Airflow中Task1无法向Task2推送数据的问题排查求助
Apache Airflow DAG任务依赖错误导致XCom拉取失败的解决方法
问题描述
在Apache Airflow中构建了一个包含两个任务的DAG:
pull_data_from_gsheet:从Google Sheet拉取数据并清洗为pandas DataFrame,转为JSON字符串后推送到XCompush_data_to_bigquery:从XCom拉取数据并写入BigQuery表
当前问题:Task2无法拉取到Task1推送的XCom数据,df_json值为None,单独运行Task1正常,Task2测试时报无数据错误。已确认XCom键匹配、DataFrame转JSON流程正常、两个PythonOperator均设置provide_context=True。
核心原因
任务依赖配置完全颠倒:代码中使用push_data_task.set_downstream(pull_data_task),该配置表示Task2是Task1的上游,即Task2会先于Task1执行。此时Task1尚未运行生成数据,Task2自然无法拉取到对应XCom内容。
解决方案
1. 修正任务依赖关系
将依赖关系调整为Task1执行完成后再执行Task2,推荐使用Airflow更直观的位运算符写法:
# 替换原有的push_data_task.set_downstream(pull_data_task) pull_data_task >> push_data_task
也可以使用set_upstream方法实现相同效果:
push_data_task.set_upstream(pull_data_task)
2. 优化XCom拉取的准确性
在Task2的xcom_pull中明确指定task_ids参数,避免多任务存在同名XCom键时的冲突:
df_json = kwargs['ti'].xcom_pull(task_ids='pull_data_from_gsheet', key='transformed_data')
3. 修正Task2的代码缩进错误
原Task2代码中,写入BigQuery的逻辑缩进错误,会在df_json为None时尝试使用未定义的df变量,修正后代码:
def push_data_to_bigquery(**kwargs): # 明确指定拉取目标任务的XCom df_json = kwargs['ti'].xcom_pull(task_ids='pull_data_from_gsheet', key='transformed_data') if df_json is not None: # 将JSON字符串转回DataFrame df = pd.read_json(df_json, orient='split') print(df) # 写入BigQuery的逻辑必须放在分支内 gbq.to_gbq(df, destination_table=bigquery_table_name, project_id=project_id, if_exists='replace') else: raise ValueError("No data found in XCom for key 'transformed_data'.")
4. 额外验证点
- 确认Task1的
xcom_push逻辑未被异常中断:可在Task1中添加异常捕获,或查看Airflow任务日志确认任务执行成功且XCom已推送 - 检查XCom存储:在Airflow UI的XCom页面,确认
pull_data_from_gsheet任务下存在transformed_data键的记录
修正后的完整任务配置代码
pull_data_task = PythonOperator( task_id='pull_data_from_gsheet', python_callable=pull_data_from_gsheet, provide_context=True, dag=dag, ) push_data_task = PythonOperator( task_id='push_data_to_bigquery', python_callable=push_data_to_bigquery, provide_context=True, dag=dag, ) # 正确的依赖关系:先执行数据拉取清洗,再执行BigQuery写入 pull_data_task >> push_data_task
内容的提问来源于stack exchange,提问作者Rajeev Pandey
相关产品推荐
相关产品推荐

