如何通过Airflow将大CSV文件流式写入BigQuery?
实现超大CSV文件流式写入BigQuery(Airflow + WRITE_APPEND模式)
针对超大CSV文件的流式写入需求,推荐优先用批量写入(效率更高、成本更低),逐行写入仅适合测试或极小批量场景。以下是具体实现方案:
一、批量流式写入(推荐)
思路
通过Airflow的PythonOperator分块读取CSV,每积累一定行数就调用BigQuery客户端批量插入,指定WRITE_APPEND模式追加数据。既避免内存溢出,又保证写入效率。
实现代码
方式1:用csv模块+insert_rows_json(轻量,无需pandas)
from airflow import DAG from airflow.operators.python import PythonOperator from google.cloud import bigquery from google.oauth2 import service_account from airflow.hooks.base import BaseHook import csv from datetime import datetime # 配置参数 PROJECT_ID = "your-project-id" DATASET_ID = "your-dataset-id" TABLE_ID = "your-table-id" CSV_FILE_PATH = "/path/to/large.csv" BATCH_SIZE = 10000 # 单批次行数,根据Worker内存调整 GCP_CONN_ID = "google_cloud_default" # Airflow中配置的GCP连接 def batch_write_to_bq(): # 从Airflow连接获取凭证 conn = BaseHook.get_connection(GCP_CONN_ID) credentials = service_account.Credentials.from_service_account_info(conn.extra_dejson) client = bigquery.Client(credentials=credentials, project=PROJECT_ID) table_ref = client.dataset(DATASET_ID).table(TABLE_ID) # 分块读取CSV并批量写入 with open(CSV_FILE_PATH, 'r', encoding='utf-8') as f: reader = csv.DictReader(f) batch = [] for row_num, row in enumerate(reader, start=1): # 预处理数据(去空格、空值处理,按需调整) processed_row = {k: v.strip() if v else None for k, v in row.items()} batch.append(processed_row) # 达到批次大小或文件末尾时执行写入 if len(batch) >= BATCH_SIZE or row_num % BATCH_SIZE == 0: errors = client.insert_rows_json(table_ref, batch) if errors: raise Exception(f"Batch insert failed at row {row_num}: {errors}") batch = [] print(f"Successfully inserted {row_num} rows") # 处理剩余的最后一批数据 if batch: errors = client.insert_rows_json(table_ref, batch) if errors: raise Exception(f"Final batch insert failed: {errors}") print(f"Total rows inserted: {row_num}") # 定义DAG default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), 'retries': 2 } with DAG( 'large_csv_to_bq_batch', default_args=default_args, schedule_interval=None, # 按需手动触发或配置定时 catchup=False ) as dag: batch_write_task = PythonOperator( task_id='batch_write_to_bigquery', python_callable=batch_write_to_bq ) batch_write_task
方式2:用Pandas分块+load_table_from_dataframe(适合复杂数据转换)
如果需要对CSV数据做类型转换、清洗,用Pandas的分块读取更便捷:
import pandas as pd def batch_write_with_pandas(): conn = BaseHook.get_connection(GCP_CONN_ID) credentials = service_account.Credentials.from_service_account_info(conn.extra_dejson) client = bigquery.Client(credentials=credentials, project=PROJECT_ID) table_ref = client.dataset(DATASET_ID).table(TABLE_ID) # 分块读取CSV为DataFrame for chunk in pd.read_csv(CSV_FILE_PATH, chunksize=BATCH_SIZE): # 数据预处理示例:字符串去空格、日期格式转换 chunk = chunk.apply(lambda x: x.str.strip() if x.dtype == 'object' else x) chunk['date_col'] = pd.to_datetime(chunk['date_col'], format='%Y-%m-%d') # 写入BigQuery,指定WRITE_APPEND模式 job = client.load_table_from_dataframe( chunk, table_ref, write_disposition=bigquery.WriteDisposition.WRITE_APPEND ) job.result() # 等待写入完成 print(f"Inserted {len(chunk)} rows")
二、逐行流式写入(不推荐)
仅适合测试或极小批量场景,因为逐行写入会产生大量API请求,速度慢且成本高:
def stream_row_by_row(): conn = BaseHook.get_connection(GCP_CONN_ID) credentials = service_account.Credentials.from_service_account_info(conn.extra_dejson) client = bigquery.Client(credentials=credentials, project=PROJECT_ID) table_ref = client.dataset(DATASET_ID).table(TABLE_ID) with open(CSV_FILE_PATH, 'r', encoding='utf-8') as f: reader = csv.DictReader(f) for row_num, row in enumerate(reader, start=1): processed_row = {k: v.strip() if v else None for k, v in row.items()} errors = client.insert_rows_json(table_ref, [processed_row]) if errors: raise Exception(f"Row {row_num} insert failed: {errors}") if row_num % 1000 == 0: print(f"Successfully inserted {row_num} rows")
三、TB级超大文件优化方案
如果CSV文件达到TB级,建议先上传到GCS,再用Airflow的CloudStorageToBigQueryOperator完成加载——这是Google内部优化的传输路径,速度更快且不占用Airflow Worker资源:
from airflow.providers.google.cloud.transfers.local_to_gcs import LocalFilesystemToGCSOperator from airflow.providers.google.cloud.operators.bigquery import CloudStorageToBigQueryOperator with DAG(...) as dag: # 上传本地CSV到GCS upload_to_gcs = LocalFilesystemToGCSOperator( task_id='upload_csv_to_gcs', src=CSV_FILE_PATH, dst='large_datasets/large.csv', bucket='your-gcs-bucket', gcp_conn_id=GCP_CONN_ID ) # 从GCS加载到BigQuery,WRITE_APPEND模式 gcs_to_bq = CloudStorageToBigQueryOperator( task_id='gcs_to_bigquery', bucket='your-gcs-bucket', source_objects=['large_datasets/large.csv'], destination_project_dataset_table=f"{PROJECT_ID}.{DATASET_ID}.{TABLE_ID}", write_disposition='WRITE_APPEND', skip_leading_rows=1, # 跳过CSV表头 gcp_conn_id=GCP_CONN_ID ) upload_to_gcs >> gcs_to_bq
关键注意事项
- Schema匹配:确保CSV列名、数据类型与BigQuery表完全一致,否则会插入失败。若表未创建,可在代码中提前定义schema并创建表。
- 错误重试:对批量写入逻辑添加重试机制(比如用
tenacity库),避免网络波动导致任务失败。 - Airflow资源:Worker节点内存需满足单批次数据处理需求,避免OOM。
内容的提问来源于stack exchange,提问作者Jooo21
相关产品推荐
相关产品推荐

