Airflow中异常如何传递至on_failure_callback及无法获取问题
解决方案
问题根因
Airflow 2.0到2.2.x版本默认将DAG级别的回调放在调度器进程执行,不会将输出写入任务实例日志,也无法直接操作对应任务的XCOM,2.1.2版本还存在小概率异常对象未正确注入上下文的缺陷。
推荐解决方法
调整调度器配置,让回调在Worker进程执行:
- 修改Airflow配置项
[scheduler] run_duration_callbacks_in_scheduler为False,Docker部署时可直接传递环境变量AIRFLOW__SCHEDULER__RUN_DURATION_CALLBACKS_IN_SCHEDULER=False到所有调度器、Worker容器,重启服务生效。
调整后回调会在执行失败任务的Worker进程中运行,上下文的异常对象可以正常读取,日志和XCOM操作都会生效。
正确的回调函数示例
import logging def known_error_dag(context): exception = context.get("exception") # 无异常信息时直接发送告警 if not exception: send_alert_email() return exception_msg = str(exception) # 匹配到指定重复错误跳过告警 if "there are duplicates" in exception_msg: return # 其余异常发送告警 send_alert_email() # 可选项:将异常存入XCOM留档 ti = context["ti"] ti.xcom_push(key="failure_exception", value=exception_msg)
无需修改全局配置的替代方案
如果不希望调整全局配置,可将异常捕获逻辑下沉到具体任务中,用try-except包裹任务逻辑,捕获到异常后判断关键词,再决定是否触发告警即可。
内容的提问来源于stack exchange,提问作者Javier Lopez Tomas
相关产品推荐
相关产品推荐

