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

如何通过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

关键注意事项

  1. Schema匹配:确保CSV列名、数据类型与BigQuery表完全一致,否则会插入失败。若表未创建,可在代码中提前定义schema并创建表。
  2. 错误重试:对批量写入逻辑添加重试机制(比如用tenacity库),避免网络波动导致任务失败。
  3. Airflow资源:Worker节点内存需满足单批次数据处理需求,避免OOM。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 22:06:03