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

Airflow中Task1无法向Task2推送数据的问题排查求助

Apache Airflow DAG任务依赖错误导致XCom拉取失败的解决方法

问题描述

在Apache Airflow中构建了一个包含两个任务的DAG:

  • pull_data_from_gsheet:从Google Sheet拉取数据并清洗为pandas DataFrame,转为JSON字符串后推送到XCom
  • push_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 21:07:21