尧图网站设计 尧图网站设计YAOTU DESIGN
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),仅供参考
返回列表