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

Airflow技术问题:如何从BigQuery表获取数据并转为列表使用?

问题原因

你的代码核心问题在于Airflow的DAG解析时机和任务运行时机不匹配:

  • run_export这个Operator是在Airflow调度器加载DAG(解析阶段)时就初始化的,这时候process_data任务还没执行,bq_data还是初始的空列表,所以访问bq_data[0]会直接抛出"列表为空"的错误。
  • process_data任务里的XCom数据只有在DAG实际运行时才会被拉取,这时候DAG已经完成解析,无法再修改run_export的参数。
解决方案

你可以通过两种方式实现动态生成CSV文件名的需求:

方案一:使用Jinja2模板直接引用XCom数据

Airflow支持用Jinja2模板在Operator参数中直接引用XCom数据,无需通过全局变量传递。注意任务ID必须在解析阶段确定,不能依赖运行时数据,所以任务ID要固定,只动态修改存储路径:

from airflow.providers.google.cloud.operators.bigquery import BigQueryGetDataOperator, BigQueryToCloudStorageOperator
from airflow.operators.python import PythonOperator
from airflow.models import DAG
from datetime import datetime

default_args = {
    'start_date': datetime(2024, 1, 1)
}

with DAG('bq_export_dag', default_args=default_args, schedule_interval=None) as dag:
    get_data = BigQueryGetDataOperator(
        task_id='get_data_from_bq',
        dataset_id='SK22',
        table_id='current_times',
        max_results=100,  # 此处传整数而非字符串
        selected_fields='current_timestamps',
    )

    def process_data_from_bq(**kwargs):
        ti = kwargs['ti']
        bq_data = ti.xcom_pull(task_ids='get_data_from_bq')
        # 若需格式化时间戳等处理,可在此操作后推回XCom
        ti.xcom_push(key='formatted_timestamp', value=bq_data[0][0])

    process_data = PythonOperator(
        task_id='process_data_from_bq',
        python_callable=process_data_from_bq,
        provide_context=True
    )

    run_export = BigQueryToCloudStorageOperator(
        task_id="save_data_on_storage",  # 任务ID固定,不可动态生成
        source_project_dataset_table="a-data-set",
        # 用Jinja2模板引用处理后的XCom数据
        destination_cloud_storage_uris=[
            "gs://europe-west1-airflow-bucket/data/test{{ ti.xcom_pull(task_ids='process_data_from_bq', key='formatted_timestamp') }}.csv"
        ],
        export_format="CSV",
        field_delimiter=",",
        print_header=False,
    )

    get_data >> process_data >> run_export

方案二:在PythonOperator中动态执行导出逻辑

如果需要更灵活的控制,可以直接在PythonOperator中使用BigQueryHook执行导出操作,在运行时拿到XCom数据后再构造导出路径:

from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook
from airflow.providers.google.cloud.operators.bigquery import BigQueryGetDataOperator
from airflow.operators.python import PythonOperator
from airflow.models import DAG
from datetime import datetime

default_args = {
    'start_date': datetime(2024, 1, 1)
}

def export_to_gcs(**kwargs):
    ti = kwargs['ti']
    # 拉取BigQuery返回的数据
    bq_data = ti.xcom_pull(task_ids='get_data_from_bq')
    if not bq_data:
        raise ValueError("未从BigQuery获取到数据")
    timestamp = bq_data[0][0]  # XCom返回格式为[(时间戳,), ...],需取嵌套元素
    # 初始化Hook执行导出
    hook = BigQueryHook(gcp_conn_id='google_cloud_default')
    hook.extract_table_to_storage(
        source_project_dataset_table="a-data-set",
        destination_cloud_storage_uris=[f"gs://europe-west1-airflow-bucket/data/test{timestamp}.csv"],
        export_format="CSV",
        field_delimiter=",",
        print_header=False
    )

with DAG('bq_export_dag', default_args=default_args, schedule_interval=None) as dag:
    get_data = BigQueryGetDataOperator(
        task_id='get_data_from_bq',
        dataset_id='SK22',
        table_id='current_times',
        max_results=100,
        selected_fields='current_timestamps',
    )

    run_export = PythonOperator(
        task_id='save_data_on_storage',
        python_callable=export_to_gcs,
        provide_context=True
    )

    get_data >> run_export
额外注意点
  • BigQueryGetDataOperator的max_results参数需传整数,而非字符串,否则可能触发异常。
  • XCom存储的数据格式为列表,每个元素是对应行的元组,因此访问首个值需要用bq_data[0][0],而非bq_data[0]。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:30:41