如何使用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:根据业务需求选择写入模式,避免数据重复或覆盖错误。
注意事项
- 确保Airflow服务账号拥有GCS读取权限和BigQuery写入权限。
- 若CSV包含特殊字符或换行,可添加
quote_character='"', allow_quoted_newlines=True参数适配。
内容的提问来源于stack exchange,提问作者user17328160
相关产品推荐
相关产品推荐

