Apache Airflow 任务日志指南:FileTaskHandler 配置、远程日志、日志分组与自定义 Handler 实战
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 围绕"每个任务实例独立产生、独立查看日志"这一核心需求,构建了以FileTaskHandler为基础的任务日志体系,并通过文件命名模板、远程日志上传、日志分组标记、日志交错合并、worker 日志实时推送等机制,让任务级日志在 Web UI 中既可完整追溯,又清晰可读。本文以 logging-tasks.rst 为主干,结合本仓库 Airflow 源码与配置模板,系统讲解任务日志的落盘路径规则、从自定义代码写日志的三种方式、::group::分组折叠、日志交错解析、worker/triggerer 实时日志服务,以及面向 trigger 场景编写自定义 FileTaskHandler 时必须掌握的四个关键属性。
任务日志的整体模型:每个 Task 一条独立日志流
Airflow 之所以能让你在 UI 中分别查看每个任务的日志,是因为它"按任务分文件"地组织日志:核心模块提供FileTaskHandler这一日志处理接口,把任务日志写入文件,并且在任务运行期间提供从 Worker 把日志实时喂给 Web UI 的机制。与此同时,Apache Airflow 社区还针对大量外部服务发布了 Provider(见 providers 概览),其中不少 Provider 提供了扩展 Airflow 日志能力的 handler,完整清单参见 providers 的日志扩展文档。
从源码看,FileTaskHandler本身继承自 Python 标准库的logging.Handler(见 file_task_handler.py),它的工作方式是:在拿到 TaskInstance 上下文后(set_context),根据命名模板计算目标日志文件路径,并用NonCachingRotatingFileHandler实际执行写入(支持max_bytes、backup_count、delay等轮转参数);读取日志时则通过模板方法_read汇总 Worker 本地文件、远程日志、Executor 日志、HTTP 服务端日志等多个来源。
值得说明的是,各服务(如 S3、GCS、WASB、HDFS、OSS)的远程 handler 大多通过把remote_base_log_folder的 scheme 与配置中的连接 ID 组合而动态选择,Airflow 默认配置文件中给出了完整的初始化逻辑(见 airflow_local_settings.py),例如以s3://开头会构造s3.task_handler、cloudwatch://对应 CloudWatch handler、gs://对应 GCS handler、wasb前缀选择 WASB handler 等。
使用 S3、GCS、WASB、HDFS 或 OSS 这类"对象存储/文件系统"型远程日志服务时,本地日志文件在上传成功后可以删除以节省磁盘,配置方式如下:
[logging] remote_logging = True remote_base_log_folder = schema://path/to/remote/log delete_local_logs = True其中delete_local_logs默认值为False,对应配置项自 2.6.0 引入(见 config.yml);其底层行为是远程 handler 初始化时读取该布尔值作为"上传后是否删除本地副本"的依据。
配置日志落盘目录与文件命名模板
对于默认 handler(FileTaskHandler),可以通过airflow.cfg中的base_log_folder指定日志文件存放目录:
[logging] base_log_folder = /your/log/folder该配置项默认值为{AIRFLOW_HOME}/logs,即默认把日志放在AIRFLOW_HOME目录下,并且文档明确要求该路径必须是绝对路径。配置模板还特别提醒(见 config.yml):如果覆盖了默认值,很可能需要同步更新[logging] dag_processor_manager_log_location与[logging] dag_processor_child_process_log_directory,因为有一批既有配置假定base_log_folder处于默认位置。
更多关于如何设置 Airflow 配置项(配置文件、环境变量、命令行)的介绍见 set-config.rst。
任务日志文件的命名遵循如下默认模式:
- 普通任务:
dag_id={dag_id}/run_id={run_id}/task_id={task_id}/attempt={try_number}.log - 动态映射(Dynamically Mapped)任务:
dag_id={dag_id}/run_id={run_id}/task_id={task_id}/map_index={map_index}/attempt={try_number}.log
即映射任务在task_id与attempt之间额外插入map_index层级,使同一映射任务的不同下游实例日志相互隔离。两条规则均可通过[logging] log_filename_template调整,其默认值正是上述两种形态的合体模板(见 config.yml):
dag_id={{ ti.dag_id }}/run_id={{ ti.run_id }}/task_id={{ ti.task_id }}/{%% if ti.map_index >= 0 %%}map_index={{ ti.map_index }}/{%% endif %%}attempt={{ try_number|default(ti.try_number) }}.log该模板使用 Jinja 语法,渲染时可用的变量包括ti(TaskInstance)、try_number等。在源码的_render_filename中,路径最终由FileTaskHandler依据该模板渲染得到(见 file_task_handler.py)。此外还可以在base_log_folder之外再提供一个远程位置(即上文remote_base_log_folder),用于存放当前日志与历史备份。
FileTaskHandler在创建新目录和新文件时,还会读取两个权限配置(见 config.yml):
file_task_handler_new_folder_permissions:新目录权限,默认0o775;file_task_handler_new_file_permissions:新文件权限,默认0o664。
这两项在使用 impersonation(以任务用户身份写日志)时尤为重要——推荐把 Airflow 用户与任务用户加入同一组并保持组可写;若确定不用 impersonation,可收紧为0o755/0o644,甚至0o700/0o600。
在业务代码中向任务日志写入内容
Airflow 使用 Python 标准logging框架输出日志,并且在任务执行的整个周期内,会把根 logger 配置为写入该任务的日志文件。绝大多数 Operator 都带有一个类型为airflow.sdk.types.Logger的log属性,凡是继承自airflow.sdk.BaseOperator的 Operator 都会自动获得该 logger,因此"Operator 内部自动写任务日志"是最常见的情形。
同时,由于任务执行期间根 logger 被指向任务日志,任何使用默认设置、且会向根 logger 传播(propagate)的标准 Python logger,其输出也会自动进入任务日志。因此从自定义代码写任务日志有三种等价途径:
- 使用
BaseOperator上自带的self.loglogger; - 使用标准
print语句输出到stdout(文档并不推荐,但某些场景确实可行); - 按 Python 惯例以模块名创建 logger 并写入。
方式三即日常最常见的写法:
import logging logger = logging.getLogger(__name__) logger.info("This is a log message")在 Airflow 3.x 的执行模型下,任务 SDK 侧通过logging_mixin等机制确保这类日志会被统一收集到任务日志流中,最终在任务日志文件中按时间顺序呈现。
日志分组:用::group::/::endgroup::折叠无关内容
版本 2.9.0 起,Airflow 支持在日志中插入分组标记,把大段日志折叠起来,让超长任务日志像 CI 流水线一样可读。该方案与 GitHub Actions 的 workflow command、Azure DevOps 的 logging command 采用兼容的分组语法,因此在 CI 中产出此类标记的工具,其输出在 Airflow UI 中可以直接获得同样的折叠体验。
在代码中按如下方式加入分组标记即可:
print("Here is some standard text.") print("::group::Non important details") print("bla") print("debug messages...") print("::endgroup::") print("Here is again some standard text.")当这些日志在 Web UI 中展示时,折叠后的内容如下(其中⯈表示"当前处于折叠态、可展开"):
[2024-03-08, 23:30:18 CET] {logging_mixin.py:188} INFO - Here is some standard text. [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} ⯈ Non important details [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} INFO - Here is again some standard text.点击"Non important details"这一行文本标签后,分组内细节展开,折叠标记变为⯆,并在组尾显示⯅⯅⯅ Log group end:
[2024-03-08, 23:30:18 CET] {logging_mixin.py:188} INFO - Here is some standard text. [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} ⯆ Non important details [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} INFO - bla [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} INFO - debug messages... [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} ⯅⯅⯅ Log group end [2024-03-08, 23:30:18 CET] {logging_mixin.py:188} INFO - Here is again some standard text.在 Airflow 新 UI 中,::group::/::endgroup::标记会被解析成可折叠的分组(group header)而非普通行。相关实现与测试位于 airflow-core/src/airflow/ui/src/pages/TaskInstance/Logs,其中 Logs.test.tsx 覆盖了"分组默认全部折叠、支持 Expand All / Collapse All、支持嵌套分组"等交互行为;服务端读取逻辑则在 log_reader.py 中也会注入::group::Log message source details/::endgroup::这类结构化分组消息,把日志来源信息也纳入可折叠区域。
日志交错:blob 型远程 handler 如何还原时间顺序
远程任务日志 handler 大体可分为两类:
- 流式 handler:如 ElasticSearch、AWS Cloudwatch、GCP operations logging(原 Stackdriver)等,无论任务处于哪个阶段、在哪个位置执行,所有消息都能以同一标识发往日志服务,通常无需从多个来源取数,也就不存在交错问题。
- blob 存储型 handler:如 S3、GCS、WASB 等。依据任务状态的不同,日志可能分散在多个位置、多个文件中(例如 Worker 日志与 Triggerer 日志分开存放、尝试重试的日志各自成文件等)。为还原完整可读的顺序日志,Airflow 需要把各来源按行取出并"交错合并",这就必须解析每一行的时间戳。
源码中的交错实现位于 file_task_handler.py:每条日志流先被解析成(timestamp, line_num, message)记录,再通过基于heapq的 K 路归并(K-way merge)按"时间戳×偏移量 + 行号"生成的排序键合并为全局有序流,同时还会对相邻重复行做去重(避免同一文件被多次读到造成重复展示)。
默认情况下,时间戳解析函数从行首取第一个以空格分隔的 token 并剥去首尾的[]后交给 pendulum 解析(见 file_task_handler.py 同级源码 file_task_handler.py)。如果你使用了自定义 formatter(时间格式非默认形态),则需要通过如下配置提供一个"自定义解析函数"的可调用路径:
[logging] interleave_timestamp_parser = path.to.my_func该 callable 接收一行日志字符串,返回一个兼容datetime.datetime的时间戳对象(配置项说明见 config.yml)。
故障排查:如何确认当前生效的 Task Handler
要快速查看当前任务日志 handler 究竟是什么,可直接运行airflow info:
$ airflow info Apache Airflow version | 2.9.0.dev0 executor | LocalExecutor task_logging_handler | airflow.utils.log.file_task_handler.FileTaskHandler sql_alchemy_conn | postgresql+psycopg://postgres:airflow@postgres/airflow dags_folder | /files/dags plugins_folder | /root/airflow/plugins base_log_folder | /root/airflow/logs remote_base_log_folder | [skipping the remaining outputs for brevity]task_logging_handler一行显示的正是当前生效的 task handler(若配置了远程 handler,则此处会显示如airflow.providers.amazon.aws.log.s3_task_handler.S3TaskHandler之类的值);base_log_folder、remote_base_log_folder则给出本地与远程日志位置。airflow info命令的输出由 info_command.py 收集并打印,可将其视为"日志配置体检报告"。
此外,还可以运行airflow config list来校验[logging]段各项配置值是否合法。当启用远程日志但 Web UI 拉取日志异常时,优先用这两条命令核对 handler 与文件夹配置。
从 Worker 与 Triggerer 实时推送日志
多数任务日志 handler 是在任务完成后才把日志发送出去。为了让用户在任务运行期间就能实时查看日志,Airflow 会启动一个轻量 HTTP 服务来 serve 本地日志,触发条件如下:
- 使用
LocalExecutor时,在airflow scheduler运行期间启动; - 使用
CeleryExecutor时,在airflow worker运行期间启动; - Triggerer 默认也服务日志,除非以
--skip-serve-logs选项启动。
该 HTTP 服务监听端口分别由[logging]段的worker_log_server_port(默认8793)与trigger_log_server_port(默认8794)指定(见 config.yml)。Webserver 与 Worker 之间的日志拉取通信使用[api]段的secret_key进行签名,因此各组件必须配置一致的secret_key,否则通信会失败。从源码看,Webserver 拉取 Worker 日志时,会基于该secret_key以HS512算法签发一个 audience 为task-instance-logs的短期 JWT,并放入请求的Authorization头(见 file_task_handler.py)。拉取超时由[api] log_fetch_timeout_sec控制,JWT 有效期内可容忍的时钟偏差由[webserver] log_request_clock_grace决定(默认 30 秒)。
底层 HTTP 服务器使用 Gunicorn(WSGI 服务器),其配置可用环境变量GUNICORN_CMD_ARGS覆盖,具体可配置项见 Gunicorn 官方 settings 文档。
编写自定义 FileTaskHandler(面向 Triggerer 的要点)
Provider 生态中已有覆盖各主流云厂商的丰富 handler,多数场景直接选用即可(清单见 providers 日志扩展文档)。只有当现有 handler 都无法满足需求、需要对接全新服务时,才建议自行实现 FileTaskHandler。这是高级话题,尤其需要注意 Trigger 日志的引入带来的一系列设计约束。
Trigger 与传统任务最大的区别在于:许多 trigger 运行在同一个进程中,且它们跑在 asyncio 事件循环之上,因此不能通过日志 handler 引入阻塞调用;同时不同 handler 行为差异很大(有的写文件、有的上传 blob、有的在消息到达时即时发网络消息、有的在独立线程中发送),Triggerer 需要某种机制来获知"应当如何使用这个 handler"。
为此,Airflow 定义了一组可设置在 handler 实例或类上的属性,用于向 Triggerer 描述该 handler 的行为。需要注意:这些参数的判定不遵循类的继承关系——因为 FileTaskHandler 的子类可能在相关特性上与父类不同,因此即使某个属性在父类上是某个值,子类也必须显式声明。四个属性如下:
trigger_should_wrap:控制该 handler 是否应被TriggerHandlerWrapper包装。当 handler 的每个实例都会创建一个文件句柄、并把收到的所有消息都写入该句柄时(典型的如 FileTaskHandler 本身),必须置为True。从源码可见FileTaskHandler类上默认trigger_should_wrap = True(见 file_task_handler.py)。trigger_should_queue:控制 Triggerer 是否应在事件循环与 handler 之间插入一个QueueListener,把 handler 里的阻塞 IO 移出事件循环,避免阻塞 asyncio 主循环。trigger_send_end_marker:控制当某个 trigger 完成时,是否向 logger 发送 END 信号。该信号用于告诉包装器关闭并移除与刚完成的 trigger 对应的那个独立文件句柄。trigger_supported:如果trigger_should_wrap与trigger_should_queue均不为True,通常就认为该 handler 不支持 trigger;但若此时 handler 显式把trigger_supported置为True,Triggerer 启动时仍会把该 handler 挂到根 logger 上,使其能处理 trigger 消息。这一属性本质上适用于"原生支持 trigger"的 handler,StackdriverTaskHandler即属此类(它能直接把消息推送到远程流式服务,天然适合 asyncio 环境)。
换言之:写文件的 handler 需要trigger_should_wrap(每个 trigger 单独一个 wrapper + 文件)、做阻塞 IO 的 handler 需要trigger_should_queue(用 QueueListener 隔离事件循环)、需要感知 trigger 结束的 handler 需要trigger_send_end_marker,而流式原生 handler 则可以只声明trigger_supported = True。从源码结构看,Triggerer 正是依据这些标记来决定对某 handler 采用"包装 + 队列"还是"直接挂接"的处理策略,从而把种类繁多的第三方 handler 统一纳入 trigger 日志体系。
远程日志的外部 UI 链接
使用远程日志时,可以配置让 Airflow Web UI 在任务日志页面显示一条指向外部日志系统 UI 的链接,点击后跳转到外部界面查看原始日志。部分外部系统要求在 Airflow 中做特定配置才能正确跳转,另一些则开箱即用;具体配置方式随 Provider 而异,例如 S3/GCS/CloudWatch 等对应的 Provider 文档中通常给出各自的外链拼接参数。
延伸阅读
- 深入配置远程日志、自定义各组件 handler、为单个 Operator/Hook/Task 定制日志 handler,见 advanced-logging-configuration.rst;
- 理解 Airflow 各组件(scheduler、worker、triggerer、webserver 等)日志的整体流向,见 logging-architecture.rst;
- 各配置项的权威默认值与说明见 config.yml 的
logging:段; - 任务日志读写核心实现见 file_task_handler.py;
- 常用远程 handler 的初始化分支见 airflow_local_settings.py。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考