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

Airflow中GCSToGCSOperator无法拉取XComs值作为参数的问题

问题

我尝试将DAG中第一个任务(task ID为get_file_list_click)返回的文件名列表,通过xcom_pull传递给第二个任务task_copy_file_list_click的GCSToGCSOperator作为source_objects参数。
运行环境为GCP Cloud Composer v1.20.12、Airflow v2.4.3,Python代码如下:

with models.DAG(
        dag_id='stg_ingestion',
        start_date=days_ago(2),
        schedule_interval='@once',
        render_template_as_native_obj=True,
    ) as dag:

    # get file list for clicks
    task_get_file_list_click = PythonOperator(
        task_id='get_file_list_click',
        provide_context=True,
        python_callable=get_file_list,
        op_args=['click'],
        dag=dag
    )

    # copy files from source to destination bucket
    task_copy_file_list_click = GCSToGCSOperator(
        task_id='copy_file_list_click',
        source_bucket=SOURCE_BUCKET,
        source_objects="{{ ti.xcom_pull(task_ids='get_file_list_click') }}",
        destination_bucket=DESTINATION_BUCKET,
        impersonation_chain=SERVICE_ACCOUNT,
        dag=dag
    )

    # load click file(s) to BigQuery
    task_load_csv_files_click = GCSToBigQueryOperator(
        task_id='load_csv_files_click',
        bucket=SOURCE_BUCKET,
        impersonation_chain=SERVICE_ACCOUNT,
        source_objects="{{ ti.xcom_pull(task_ids='get_file_list_click') }}",    
        compression='GZIP',
        destination_project_dataset_table=f"{DATASET_NAME}.dcm_click",
        schema_object_bucket=SCHEMA_BUCKET,
        schema_object="dags/resources/json/dcm_click.json",        
        write_disposition='WRITE_TRUNCATE',
        skip_leading_rows=1,
        trigger_rule='all_success',
        dag=dag,
    )

第二个任务报错:"NoneType object is not subscriptable",但第三个任务task_load_csv_files_click的GCSToBigQueryOperator却能成功通过xcom_pull拉取文件名列表作为source_objects参数运行。

解决建议
  • 核对XCom推送结果:先在Airflow UI的任务实例详情页查看get_file_list_click的XCom数据,确认是否返回了预期的非空文件名列表。如果XCom值为None,说明get_file_list函数未正确返回数据,或PythonOperator未完成推送。
  • 添加模板默认值:修改GCSToGCSOperator的source_objects模板,指定默认空列表,避免None传入引发错误:
    source_objects="{{ ti.xcom_pull(task_ids='get_file_list_click', default=[]) }}"
    
  • 显式声明XCom推送:虽然PythonOperator默认开启do_xcom_push=True,但显式声明可避免环境或函数逻辑导致的推送异常:
    task_get_file_list_click = PythonOperator(
        task_id='get_file_list_click',
        provide_context=True,
        python_callable=get_file_list,
        op_args=['click'],
        do_xcom_push=True,
        dag=dag
    )
    
  • 检查任务依赖:显式添加任务依赖关系,确保task_copy_file_list_click在get_file_list_click成功执行后才启动:
    task_get_file_list_click >> task_copy_file_list_click >> task_load_csv_files_click
    
  • 验证Operator参数处理差异:Airflow 2.4.3中,GCSToGCSOperator与GCSToBigQueryOperator对source_objects的空值容错逻辑不同,前者会直接遍历传入值,后者可能对空值做了兼容处理,因此必须确保传入的是有效列表。

内容的提问来源于stack exchange,提问作者rjk90

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 16:05:25