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

Airflow中GCSToBigQueryOperator传递XCom至source_objects参数失败

问题

我有一个包含两个任务的Airflow DAG:

  1. 第一个任务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
  1. 第二个任务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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 08:14:53