You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Google Cloud Composer的Airflow DAG回调中如何使用日志记录器?

在Cloud Composer的Airflow回调中记录日志的解决方案

我们在Cloud Composer(版本composer-2.0.32-airflow-2.3.4)环境中运行Airflow时,会遇到DAG回调(如on_failure_callback)中日志无法在常规位置找到的问题——即使确认回调已执行,日志也不会出现在worker、scheduler或dag-processor-manager的日志里。以下是可行的解决办法:

方法1:利用回调上下文里的任务实例日志器

Airflow的回调函数会接收context参数,其中包含当前任务实例(task_instance)的引用,直接使用它的log属性记录日志,日志会自动关联到对应任务实例的日志中,在Cloud Composer的任务日志页面就能找到:

def on_failure_callback(context):
    # 从上下文获取任务实例
    task_instance = context['task_instance']
    # 使用任务实例的日志器记录
    task_instance.log.info("回调执行失败,这是日志内容")
    # 也可以直接使用上下文里的logger对象
    context['logger'].error("另一种回调日志记录方式")

方法2:指定Airflow任务专用日志命名空间

如果偏好使用Python标准logging模块,在回调里使用airflow.task命名空间的日志器,它会自动适配Airflow的日志收集规则,日志会被归集到对应任务的日志中:

import logging

def on_failure_callback(context):
    # 使用airflow.task日志器
    log = logging.getLogger("airflow.task")
    log.setLevel(logging.INFO)
    log.info("使用airflow.task日志器的回调日志")

方法3:确认Scheduler日志级别配置

on_failure_callback默认在Scheduler进程中执行,如果日志级别设置过高(比如只捕获ERROR),INFO级别的日志会被过滤。可以在Cloud Composer环境的配置中,将Scheduler的日志级别调整为INFO:

  • 进入Cloud Composer环境详情页
  • 切换到"配置"标签,找到"Airflow scheduler日志级别"
  • 设置为INFO并保存更新

这样Scheduler进程中产生的回调日志就会被正常捕获,可以在Cloud Logging的cloud-composer/environments/[你的环境名]/airflow-scheduler日志流中查看。

内容的提问来源于stack exchange,提问作者Jodiug

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 01:46:18