在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
相关产品推荐
相关产品推荐

