Airflow中GCSToBigQueryOperator传递XCom至source_objects参数失败
我有一个包含两个任务的Airflow DAG:
- 第一个任务
list_gcs_files为PythonOperator,定义如下:
list_gcs_files = PythonOperator( task_id='list_gcs_files', python_callable=get_file_list, op_kwargs={'execution_date': '{{ execution_date }}'}, provide_context=True, )
其调用的Python函数get_file_list通过GCSHook获取指定前缀的GCS文件列表并返回:
def get_file_list(execution_date, **kwargs): bucket_name = 'some-dummy-bucket' execution_date = pendulum.parse(execution_date).date() prefix = f"customers/{execution_date.strftime('%Y-%m')}-{execution_date.day}/" hook = GCSHook(google_cloud_storage_conn_id='google_cloud_default') blobs = hook.list(bucket_name=bucket_name, prefix=prefix) files = [str(blob) for blob in blobs] return files
- 第二个任务
load_to_bronze为GCSToBigQueryOperator,需将数据加载至bronze层,定义如下:
load_to_bronze = GCSToBigQueryOperator( task_id='load_to_bronze', bucket='some-dummy-bucket', source_objects="{{ task_instance.xcom_pull(task_ids='list_gcs_files') }}", destination_project_dataset_table='bronze.customers', schema_fields=[ # SOME SCHEMA IS DEFINED HERE ], source_format='CSV', skip_leading_rows=1, write_disposition='WRITE_TRUNCATE', )
第一个任务执行成功,其XCom返回值为:
['customers/2022-08-2/2022-08-1_part2__customers.csv', 'customers/2022-08-2/2022-08-2_part1__customers.csv']
但第二个任务执行失败,报错信息为:
Source URI must not contain the ',' character: gs://some-dummy-bucket/['customers/2022-08-2/2022-08-1_part2__customers.csv', 'customers/2022-08-2/2022-08-2_part1__customers.csv']
我曾尝试将返回的XCom数据改为逗号拼接的字符串,再在第二个任务中拆分回列表,但无效;手动硬编码相同的文件列表作为source_objects参数则可正常执行。请问为何通过XCom传递列表时会失败?
核心原因
用Jinja模板{{ task_instance.xcom_pull(...) }}获取XCom列表时,Jinja会把Python列表序列化为字符串形式(比如["file1.csv","file2.csv"]),而非直接传递列表对象。
GCSToBigQueryOperator的source_objects参数支持两种输入:单个文件路径字符串,或多个文件路径组成的Python列表。但如果传入序列化后的列表字符串,Operator会将其视为单个完整的文件路径,拼接bucket后就会生成包含方括号、逗号的非法URI,触发报错。
解决方案
方案1:使用XComArg(推荐,Airflow 2.x+)
Airflow 2.x引入的XComArg可以直接传递Python对象,避免Jinja序列化问题。修改第二个任务的source_objects参数:
from airflow.models.xcom_arg import XComArg # 省略其他代码 load_to_bronze = GCSToBigQueryOperator( task_id='load_to_bronze', bucket='some-dummy-bucket', source_objects=XComArg(list_gcs_files), # 直接引用上游任务的XCom返回值 destination_project_dataset_table='bronze.customers', schema_fields=[ # 你的Schema定义 ], source_format='CSV', skip_leading_rows=1, write_disposition='WRITE_TRUNCATE', )
这样source_objects会直接拿到上游返回的Python列表,和手动硬编码效果一致。
方案2:Jinja模板处理(Airflow 1.x兼容)
如果使用Airflow 1.x,可通过Jinja过滤器将序列化的字符串转回列表:
source_objects="{{ task_instance.xcom_pull(task_ids='list_gcs_files') | replace('[','') | replace(']','') | replace(\"'\",'') | split(', ') }}"
注意:这种方式存在局限性,若文件名包含引号、逗号等特殊字符会出错,仅适合简单场景。
方案3:上游任务直接触发下游逻辑(不推荐)
在get_file_list函数中直接调用GCSToBigQueryOperator的执行逻辑,或使用BranchPythonOperator,但会增加任务耦合度,不利于维护。
内容的提问来源于stack exchange,提问作者Andrii

