免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

使用 Apache Airflow Amazon Provider 将任务日志写入 Amazon CloudWatch

使用 Apache Airflow Amazon Provider 将任务日志写入 Amazon CloudWatch 使用 Apache Airflow Amazon Provider 将任务日志写入 Amazon CloudWatch【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow将 Airflow 任务日志接入 Amazon CloudWatch Logs是构建可审计、可检索、可持续保留的日志体系的标准做法。本文以 Apache Airflow 的 Amazon Providerapache-airflow-providers-amazon为核心完整讲解如何通过airflow.cfg配置cloudwatch://远程日志后端、如何建立可读写的 Airflow AWS 连接并结合仓库源码深入剖析CloudwatchTaskHandler、CloudWatchRemoteLogIO与AwsLogsHook的底层实现与读取链路。读完本文你将能够在自己的 Airflow 实例上启用 CloudWatch 远程日志理解日志流log stream的命名规则与实时写入机制并掌握常见问题的排查方法。一、功能概览与适用场景Airflow 默认将任务日志写入调度器或 Worker 本地磁盘[logging] base_log_folder默认${AIRFLOW_HOME}/logs本地存储存在容量受限、难以跨实例聚合、无法统一检索的问题。Amazon Provider 提供的 CloudWatch 远程日志能力让每个任务实例Task Instance的日志被实时写入 AWS CloudWatch Logs 的指定 Log Group并在 Airflow Web UI 查看任务日志时自动从 CloudWatch 回读展示。这一功能的核心特征是实时流式写入配合watchtower库与后台队列线程而非传统的任务结束后批量上传。从源码实现看cloudwatch_task_handler.py日志写入与回读均由统一的 IO 抽象层完成其关键类为CloudwatchTaskHandler继承自 Airflow 核心的FileTaskHandlerfile_task_handler.py负责把本地日志文件句柄桥接到 CloudWatch并在任务结束时清理本地副本CloudWatchRemoteLogIO真正的读写执行者负责构建watchtower.CloudWatchLogHandler、调用AwsLogsHook读取日志事件AwsLogsHook对boto3.client(logs)的薄封装提供分页拉取日志事件的能力logs.py。对应的官方文档原文位于 cloud-watch-task-handlers.rst本目录下还包含同为远程日志方案的 s3-task-handler.rst二者配置模式一致只是存储后端不同。二、前置条件必须先配置好 Airflow AWS 连接原文档明确强调Remote logging to Amazon Cloudwatch uses an existing Airflow connection to read or write logs. If you dont have a connection properly setup, this process will fail.远程日志使用已有的 Airflow 连接来读写日志如果连接未正确配置该流程将失败。因此在启用 CloudWatch 日志前必须先在 Airflow 中创建一个 AWS 连接Connection其Conn Id需要与配置文件中的remote_log_conn_id一致。该连接对应的 IAM 身份至少需要具备以下权限写入logs:CreateLogStream、logs:PutLogEvents向指定 Log Group / Log Stream 写入日志事件读取logs:GetLogEventsWeb UI 回读日志时使用见AwsLogsHook.get_log_events实现。连接凭证可通过 Airflow Web UIAdmin → Connections、airflow connections add命令或环境变量注入等方式创建。Amazon Provider 的连接创建方式与其他 AWS 服务一致底层由AwsBaseHook统一处理凭证解析base_aws.py。三、核心配置airflow.cfg 中的三个关键项在原文档给出的示例基础上完整的[logging]配置如下配置项元信息可对照 config.yml 中的定义[logging] # 是否启用远程日志默认 False remote_logging True # CloudWatch 远程日志地址必须以 cloudwatch:// 开头 # 后面紧跟 Log Group 的 ARN remote_base_log_folder cloudwatch://arn:aws:logs:region name:account id:log-group:group name # 提供读写 CloudWatch 权限的 Airflow 连接 ID remote_log_conn_id MyCloudwatchConn3.1remote_logging总开关置为True后 Airflow 的日志配置模板才会进入远程日志分支。该值由 airflow_local_settings.py 中的conf.getboolean(logging, remote_logging)读取。默认值为False。3.2remote_base_log_folderCloudWatch 的地址格式这个值同时承担了选择日志处理器和定位 Log Group的双重职责前缀cloudwatch://用于让 Airflow 判定应加载 CloudWatch 处理器。在 airflow_local_settings.py 中配置模板依次匹配s3://、cloudwatch://、gs://、wasb、stackdriver://、oss://、hdfs://等前缀命中cloudwatch://后导入CloudWatchRemoteLogIO并实例化前缀之后紧跟的arn:aws:logs:region:account-id:log-group:group-name是 CloudWatch Log Group 的 ARN。源码中通过urlsplit(remote_base_log_folder)取出netloc path作为log_group_arn再按:拆分下标[3]为region_name区域下标[6]为log_groupLog Group 名称。例如arn:aws:logs:us-east-1:123456789012:log-group:my-airflow-logs会被解析为区域us-east-1、Log Group 名my-airflow-logs。注意配置中的region name、account id、group name必须替换为真实值且该 Log Group 应预先在 AWS 控制台创建或由具有建组权限的身份在首次写入时自动创建取决于 IAM 策略。3.3remote_log_conn_id指定用于 CloudWatch 读写的 Airflow 连接 ID。文档明确指出上述示例配置下 Airflow 将尝试使用AwsLogsHook(MyCloudwatchConn)。这一点在源码中得到印证cloudwatch_task_handler.py 中CloudWatchRemoteLogIO.hook构造AwsLogsHook(aws_conn_idconf.get(logging, remote_log_conn_id), region_nameself.region_name)——连接 ID 直接取自该配置项region 则从 Log Group ARN 中解析二者共同决定最终连接哪个区域的 CloudWatch Logs。3.4 相关辅助配置项除文档给出的三项外以下[logging]配置项与 CloudWatch 日志行为直接相关配置项默认值作用delete_local_logsFalse日志写入远程后是否删除本地日志副本。CloudWatch 场景下由CloudWatchRemoteLogIO.delete_local_copy使用见 config.ymlremote_task_handler_kwargsJSON以 JSON 字典形式向远程任务处理器__init__传入额外参数优先级高于 Airflow 配置值例如{delete_local_copy: true}可覆盖delete_local_logsFalse见 config.ymlbase_log_folder{AIRFLOW_HOME}/logs本地日志的基础目录远程 IO 的base_log_folder取自该值且必须是绝对路径cloudwatch_task_handler_json_serializer[aws]段无通过conf.getimport导入自定义 JSON 序列化函数用于控制写入 CloudWatch 时事件 JSON 的序列化方式见 cloudwatch_task_handler.py四、配置生效机制airflow_local_settings 的分发逻辑配置并非直接生效而是经由 Airflow 的日志配置模板完成处理器装配。在 airflow_local_settings.py 中当REMOTE_LOGGING即remote_logging True时读取remote_base_log_folder必须配置否则conf.get_mandatory_value直接报错读取remote_task_handler_kwargs并强制校验其为 JSON 对象dict否则抛出ValueError: logging/remote_task_handler_kwargs must be a JSON object (a python dict), ...按前缀匹配分发命中cloudwatch://后将base_log_folder、remote_base、delete_local_copy、log_group_arn由urlsplit解析出的netloc path以及用户自定义 kwargs 一起传入CloudWatchRemoteLogIO构造器最终将 IO 实例赋值给REMOTE_TASK_LOG供task日志处理器使用。这段分发逻辑意味着CloudWatch 与 S3、GCS、WASB、Stackdriver、OSS、HDFS 等远程后端共用同一套[logging]配置入口仅凭remote_base_log_folder前缀即可切换后端运维迁移成本低。若所有已知前缀均未命中且未配置 Elasticsearch/OpenSearch则会抛出AirflowException提示检查remote_base_log_folder配置。五、底层实现写入链路与读取链路5.1 写入链路实时流式上传与任务结束后一次性上传不同CloudWatch 日志采用实时流式写入。关键逻辑集中在 cloudwatch_task_handler.py 的_build_handler()与handler属性def _build_handler(self) - watchtower.CloudWatchLogHandler: _json_serialize conf.getimport(aws, cloudwatch_task_handler_json_serializer, fallbackNone) return watchtower.CloudWatchLogHandler( log_group_nameself.log_group, log_stream_nameself.log_stream_name, use_queuesTrue, boto3_clientself.hook.get_conn(), json_serialize_default_json_serialize or json_serialize_legacy, )要点底层使用第三方库watchtower的CloudWatchLogHandler并设置use_queuesTrue即通过队列 后台线程异步批量写入避免阻塞任务执行boto3_client来自AwsLogsHook.get_conn()即第二节提到的连接默认采用json_serialize_legacy序列化器datetime对象序列化为 ISO 格式其余非 JSON 可序列化对象序列化为null复刻 watchtower 2.0.1 行为也可通过[aws] cloudwatch_task_handler_json_serializer配置导入自定义函数替换为兼容结构日志structlogprocessors中注册的处理器会在每条日志记录上动态设置handler.log_stream_name将日志路径中的:替换为_因为 CloudWatch Log Stream 名不允许冒号。5.2 日志流命名冒号替换规则每个任务实例对应 CloudWatch 中的一个 Log Stream其命名源自FileTaskHandler._render_filename(ti, try_number)渲染出的本地日志相对路径。由于 CloudWatch Log Stream 名称不允许包含:字符源码在 cloudwatch_task_handler.py 中统一将冒号替换为下划线def _render_filename(self, ti, try_number): # Replace unsupported log group name characters return super()._render_filename(ti, try_number).replace(:, _)同理读取阶段get_cloudwatch_logs也会对目标 stream 名执行stream_name.replace(:, _)保证读写两侧命名一致。5.3 读取链路分页拉取日志事件Web UI 查看任务日志时CloudWatchRemoteLogIO.read()/stream()调用AwsLogsHook.get_log_events()从 CloudWatch 拉取事件。其分页逻辑在 logs.py 中实现几个值得注意的细节循环调用self.conn.get_log_events(...)每次携带nextToken继续读取startFromHeadTrue表示从头读取终止条件当continuation_token.value response[nextForwardToken]即 nextForwardToken 不再前进说明已读到流末尾或连续 3 次空响应NUM_CONSECUTIVE_EMPTY_RESPONSE_EXIT_THRESHOLD 3时退出。注释中说明这是参考 AWS 团队建议的折中方案——纯依赖 nextForwardToken 判读末尾可能耗时约 20 秒因此增加空响应计数加速退出参见 PR apache/airflow#20814 的说明结束时间缓冲若任务实例存在end_date读取时将其加上 30 秒作为end_time缓冲避免遗漏任务结束瞬间写入的日志优雅降级当目标 Log Stream 不存在ResourceNotFoundException时不抛 500 错误而是生成一条提示性日志事件No log stream found in CloudWatch (log_group..., log_stream...). The task may have logged to stdout only, not produced any logs yet, or remote logging may be misconfigured.帮助用户区分没日志与配置错误。5.4 任务结束清理CloudwatchTaskHandler.close()cloudwatch_task_handler.py负责在任务结束时 flush 挂起事件并清理本地副本通过closed标志防止重复上传Airflow 的logging.shutdown可能触发多次 close关闭 IO 实际使用的 handler而非set_context时缓存的引用避免 dictConfig 重建 handler 后关闭过期句柄导致后台线程泄漏调用io.upload()flush 一次后若delete_local_copyTrue则删除本地日志目录若待删除路径越出base_log_folder仅记录 warning 并跳过删除避免误删。六、验证与测试仓库为 CloudWatch 日志处理器提供了完整的单元测试见 test_cloudwatch_task_handler.py。测试使用motomock_aws模拟 AWS 服务无需真实 AWS 账号即可验证读写闭环主要覆盖通过conf_vars注入[logging]配置后构造 handler任务日志写入后被正确推送至模拟的 CloudWatch Log Group / Log Stream从 CloudWatch 回读日志事件并按[timestamp] message格式渲染_event_to_str将毫秒时间戳格式化为%Y-%m-%dT%H:%M:%SZ的 UTC 时间串本地日志删除逻辑含路径越界保护与ResourceNotFoundException场景下的提示消息生成。测试夹具中还特别说明watchtower.CloudWatchLogHandler会派生队列工作线程若不显式 close 会在测试间泄漏 handler 并阻塞logging.config.dictConfig因此测试统一在mock_aws仍生效时清理 handler——这一细节对读者自行编写集成测试同样具有参考价值。在真实环境中验证的最简路径配置完成后运行任意 DAG 任务然后在 AWS 控制台 CloudWatch Logs 中查看对应 Log Group 是否出现以dag_id/task_id/...命名的 Log Stream任务运行中即可看到实时日志说明写入链路正常随后在 Airflow Web UI 打开该任务实例的 Log 页签若能看到同样的日志内容则说明读取链路正常。七、常见问题与排查思路现象可能原因与排查方向启动 Airflow 报错 Incorrect remote log configuration...remote_base_log_folder未以cloudwatch://开头或未匹配任何已知后端前缀检查[logging]配置见 airflow_local_settings.py报错 remote_task_handler_kwargs must be a JSON object该配置项必须为合法的 JSON 字典例如{delete_local_copy: true}不能用普通字符串任务日志未出现在 CloudWatch连接 ID 是否与remote_log_conn_id一致、凭证 IAM 是否具备logs:PutLogEvents/logs:CreateLogStream权限、region 与 Log Group ARN 是否匹配Web UI 显示 No log stream found in CloudWatch...该任务可能只向 stdout 输出而未走远程处理器或尚未产生日志此提示本身就是 CloudWatch 后端给出的排查线索本地磁盘日志堆积确认delete_local_logsTrue或通过remote_task_handler_kwargs传入{delete_local_copy: true}覆盖日志中:字符异常CloudWatch Log Stream 不支持冒号Airflow 已自动将:替换为_无需手动干预八、小结Apache Airflow 的 Amazon Provider 通过cloudwatch://前缀的remote_base_log_folder、remote_log_conn_id与remote_logging三个配置项即可接入 CloudWatch Logs实现任务日志的实时远程写入与 Web UI 回读。其底层由 CloudwatchTaskHandler继承FileTaskHandler、CloudWatchRemoteLogIO基于 watchtower 实时流式写入与 AwsLogsHook分页读取日志事件协作完成且与 S3、GCS 等远程后端共用同一套配置分发框架airflow_local_settings.py运维切换成本低。配置前务必先准备好具备 CloudWatch Logs 读写权限的 Airflow AWS 连接这是整个流程能够运转的前提。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表