Airflow DAG执行报错:XComArg结果未找到,求排查方案
Airflow DAG XComNotFound 问题排查请求
DAG 核心代码片段
default_args = { 'owner': 'airflow', } @dag(default_args=default_args, start_date=datetime.datetime(2021, 1, 1), schedule_interval=None, tags=['mimir']) def data_comparism_psql(): @task(task_id="pull_release_data_from_blob") # 省略该任务实现 ... @task(task_id="deviation_calculation_metrik_calls_mean_execution_time_std_dev") def compare_release_files(release_data_export: List[Dict], querytype: Querytype): # 省略任务逻辑,返回包含"significant_differences"键的字典 ... # 任务调用逻辑 release_data_calls = get_files_from_blob(querytype=Querytype.CALLS) report_calls = compare_release_files(release_data_calls, querytype=Querytype.CALLS)["significant_differences"] report_calls_std_dev = compare_release_files(release_data_calls, querytype=Querytype.STD_DEV)["significant_differences"] release_data_mean_exec_time = get_files_from_blob(querytype=Querytype.MEAN_EXEC_TIME) report_mean_exec_time = compare_release_files(release_data_mean_exec_time, querytype=Querytype.MEAN_EXEC_TIME)["significant_differences"] report_mean_exec_time_std_dev = compare_release_files(release_data_mean_exec_time, querytype=Querytype.STD_DEV)["significant_differences"] create_release_report([report_calls, report_mean_exec_time, report_calls_std_dev, report_mean_exec_time_std_dev], "Queries") data_comparism_dag = data_comparism_psql()
报错信息
File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 75, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/xcom_arg.py", line 342, in resolve raise XComNotFound(ti.dag_id, task_id, self.key) airflow.exceptions.XComNotFound: XComArg result from deviation_calculation_metrik_calls_mean_execution_time_std_dev at data_comparism_psql with key="significant_differences" is not found!
已完成的排查动作
- 在Airflow UI中确认所有前置任务均成功执行
- 检查
deviation_calculation_metrik_calls_mean_execution_time_std_dev任务的XCom,确认存在significant_differences键及对应值 - 尝试调整任务依赖关系,问题仍未解决
排查建议
- 修复重复task_id问题:多次调用
compare_release_files但共用同一个task_id,Airflow会将这些调用视为同一个任务实例,后执行的任务会覆盖先执行的XCom数据。给每个调用分配唯一task_id:report_calls = compare_release_files.override(task_id="compare_calls")(release_data_calls, querytype=Querytype.CALLS)["significant_differences"] report_calls_std_dev = compare_release_files.override(task_id="compare_calls_std_dev")(release_data_calls, querytype=Querytype.STD_DEV)["significant_differences"] # 其余调用同理修改task_id - 验证任务返回结构:在
compare_release_files任务中添加日志,确认返回的字典确实包含significant_differences键,且值不为空。 - 检查Airflow版本bug:你使用的Python3.7搭配的Airflow版本可能存在XCom嵌套取值的兼容性问题,尝试升级到对应分支的稳定小版本。
- 核对任务实例日志:查看
deviation_calculation_metrik_calls_mean_execution_time_std_dev所有执行实例的日志,确认每次执行都成功写入XCom,无写入失败情况。 - 重新触发DAG:清空现有任务实例,重新触发整个DAG,观察每个
compare_release_files实例的执行顺序和XCom写入状态,排查是否存在数据覆盖或丢失场景。
内容的提问来源于stack exchange,提问作者AI Humanizer
相关产品推荐
相关产品推荐

