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

