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

Apache Airflow:如何在Python DAG中直接将BigQuery结果写入GCS(无中间表)

问题需求

现有基于Airflow的Python DAG代码(原代码依赖中间表存储BigQuery查询结果,再导出到GCS),需修改为直接将BigQuery查询结果以CSV格式写入GCS Bucket,且CSV文件名必须包含查询中用到的max(week_start_date)值。原代码如下:

# Query results written to intermediate table
BQ_Output =  BigQueryOperator(
    task_id='BQ_Output',
    write_disposition='WRITE_TRUNCATE',
    use_legacy_sql=False,
    allow_large_results=True, 
    sql="""
    CREATE OR REPLACE TABLE `schema.result_table` as (
    SELECT * except (store_name, state, week_start_date, season) FROM `schema.result_table` where week_start_date = (SELECT MAX(week_start_date) FROM `schema.result_table`))
    """, 
    params=var_config,
    dag=dag)
  
  
# Results Table to GCS Bucket  
Bq_GCS = BigQueryToCloudStorageOperator(
    task_id = 'Bq_GCS',
    source_project_dataset_table = "schema.result_table",
    destination_cloud_storage_uris = "gs://bucket_path/output_<max_date>_file.csv",
    export_format = 'CSV',
    field_delimiter = ',',
    dag = dag
)
解决方案

可以通过Airflow的PythonOperator结合BigQueryHook获取最大日期,再用BigQueryToCloudStorageOperator直接执行查询并导出,无需中间表。完整代码如下:

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

def get_max_date(**context):
    # 初始化BigQuery Hook
    bq_hook = BigQueryHook(use_legacy_sql=False)
    # 查询最大日期
    max_date_query = "SELECT MAX(week_start_date) as max_date FROM `schema.result_table`"
    # 执行查询并获取结果
    result = bq_hook.get_first(max_date_query)
    max_date_str = result[0].strftime("%Y-%m-%d")  # 格式化为YYYY-MM-DD,可按需调整
    # 将日期存入XCom供后续任务使用
    context['ti'].xcom_push(key='max_date', value=max_date_str)

with DAG(
    dag_id="bq_direct_export_to_gcs",
    schedule_interval="@daily",
    start_date=datetime(2024, 1, 1),
    catchup=False
) as dag:

    # 任务1:获取最大日期
    fetch_max_date = PythonOperator(
        task_id="fetch_max_date",
        python_callable=get_max_date,
        provide_context=True
    )

    # 任务2:直接将查询结果导出到GCS
    bq_to_gcs = BigQueryToCloudStorageOperator(
        task_id="bq_to_gcs",
        sql="""
        SELECT * except (store_name, state, week_start_date, season) 
        FROM `schema.result_table` 
        where week_start_date = (SELECT MAX(week_start_date) FROM `schema.result_table`)
        """,
        destination_cloud_storage_uris="gs://bucket_path/output_{}_file.csv".format("{{ ti.xcom_pull(key='max_date', task_ids='fetch_max_date') }}"),
        export_format='CSV',
        field_delimiter=',',
        use_legacy_sql=False,
        params=var_config  # 保留原有的参数配置
    )

    # 设置任务依赖
    fetch_max_date >> bq_to_gcs

关键说明

  • 获取最大日期:通过BigQueryHook直接执行查询获取最大日期,并存入XCom供后续任务调用,避免依赖中间表
  • 动态文件名:使用Airflow模板语法{{ ti.xcom_pull(...) }}获取之前任务的日期值,拼接成目标GCS路径
  • 直接导出:BigQueryToCloudStorageOperator支持通过sql参数直接传入查询语句,无需先写入中间表,减少不必要的数据存储环节

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 06:06:07