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
相关产品推荐
相关产品推荐

