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

如何使用Airflow将GCS Bucket中的文件加载至BigQuery

从GCS Bucket加载文件到BigQuery的Airflow实现

你已经完成了从BigQuery导出数据到GCS的流程,现在可以通过Airflow的CloudStorageToBigQueryOperator实现反向操作——将GCS中的CSV文件加载到BigQuery。以下是整合原有导出逻辑与新增导入任务的完整代码:

from airflow.contrib.operators.bigquery_operator import BigQueryOperator
from airflow.contrib.operators.bigquery_to_gcs import BigQueryToCloudStorageOperator
from airflow.contrib.operators.gcs_to_bigquery import CloudStorageToBigQueryOperator

# 原有导出逻辑
bq_recent_questions_query = BigQueryOperator(
    task_id='bq_recent_questions_query',
    sql="""
    SELECT owner_display_name, title, view_count
    FROM `bigquery-public-data.stackoverflow.posts_questions`
    WHERE creation_date < CAST('{max_date}' AS TIMESTAMP)
        AND creation_date >= CAST('{min_date}' AS TIMESTAMP)
    ORDER BY view_count DESC
    LIMIT 100
    """.format(max_date=max_query_date, min_date=min_query_date),
    use_legacy_sql=False,
    destination_dataset_table=bq_recent_questions_table_id)

export_questions_to_gcs = BigQueryToCloudStorageOperator(
    task_id='export_recent_questions_to_gcs',
    source_project_dataset_table=bq_recent_questions_table_id,
    destination_cloud_storage_uris=[output_file],
    export_format='CSV')

# 新增:从GCS加载到BigQuery的任务
load_gcs_to_bq = CloudStorageToBigQueryOperator(
    task_id='load_recent_questions_from_gcs_to_bq',
    # GCS文件路径,需去掉gs://前缀
    source_objects=[output_file.split('gs://')[1]],
    # 目标BigQuery表完整ID,格式:项目ID.数据集ID.表名
    destination_project_dataset_table='your-project.your-dataset.target_table',
    # 跳过CSV表头(导出的文件包含表头)
    skip_leading_rows=1,
    # 自动检测表结构,也可手动指定schema_fields
    autodetect=True,
    # 写入模式:WRITE_TRUNCATE覆盖原有数据,WRITE_APPEND追加,WRITE_EMPTY仅空表写入
    write_disposition='WRITE_TRUNCATE',
    source_format='CSV',
    field_delimiter=',')

# 设置任务依赖:查询→导出到GCS→加载回BigQuery
bq_recent_questions_query >> export_questions_to_gcs >> load_gcs_to_bq

关键参数说明

  • source_objects:GCS文件的相对路径,例如导出路径是gs://your-bucket/stackoverflow/questions.csv,这里填your-bucket/stackoverflow/questions.csv;若有多个文件,可使用通配符your-bucket/stackoverflow/questions_*.csv。
  • destination_project_dataset_table:替换为你实际要写入的BigQuery表ID。
  • schema_fields:若不需要自动检测结构,可手动定义字段类型,示例:
    schema_fields=[
        {'name': 'owner_display_name', 'type': 'STRING'},
        {'name': 'title', 'type': 'STRING'},
        {'name': 'view_count', 'type': 'INTEGER'}
    ]
    
  • write_disposition:根据业务需求选择写入模式,避免数据重复或覆盖错误。

注意事项

  1. 确保Airflow服务账号拥有GCS读取权限和BigQuery写入权限。
  2. 若CSV包含特殊字符或换行,可添加quote_character='"', allow_quoted_newlines=True参数适配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 08:40:31