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

如何从GCS无重复加载数据到BigQuery?求性能成本最优方案

问题

希望从Google Cloud Storage(Bucket)中获取CSV文件,加载到BigQuery表中并避免重复数据,要求实现性能与成本最优。当前代码如下:

def load_data_in_BQT():
    job_config = bigquery.LoadJobConfig(
        schema=[
            bigquery.SchemaField("id", "INTEGER"),
            bigquery.SchemaField("name", "STRING"),
        ],
        
        # The source format defaults to CSV, so the line below is optional.
        source_format=bigquery.SourceFormat.CSV,
        skip_leading_rows=1,
        autodetect=True,
        write_disposition=bigquery.WriteDisposition.WRITE_APPEND, # (Addition of the data (possibility of having duplications)
        # write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE, # (Formatting of the table and insertion of the new data (Loss of the old data))
    )

    uri = "gs://mybucket/myfolder/myfile.csv"
  
    load_job = self.client.load_table_from_uri(
        uri, self.table_ref["object"], job_config=job_config, 
    )

原本思路是将CSV和BigQuery表数据都转为Pandas DataFrame后去重,再用WRITE_TRUNCATE重新插入,但海量数据下存在弊端,求更好方案。


最优解决方案

以下方案均基于BigQuery原生能力实现,无需拉取数据到本地,兼顾性能与成本:

方案1:临时表+MERGE语句(海量数据首选)

完全在BigQuery内部完成数据加载与去重,避免本地内存瓶颈,是处理大规模数据的最优方式。

步骤逻辑

  1. 将CSV加载到临时表(用WRITE_TRUNCATE保证临时表每次都是最新的CSV数据)
  2. 执行MERGE语句,根据唯一键(比如id)合并临时表与目标表:不存在的插入,存在的可选择忽略或更新

代码示例

def load_data_without_duplicates():
    # 1. 配置临时表加载任务
    temp_table_ref = self.client.dataset("your_dataset").table("temp_csv_data")
    job_config = bigquery.LoadJobConfig(
        schema=[
            bigquery.SchemaField("id", "INTEGER"),
            bigquery.SchemaField("name", "STRING"),
        ],
        source_format=bigquery.SourceFormat.CSV,
        skip_leading_rows=1,
        write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE,  # 每次覆盖临时表
    )

    uri = "gs://mybucket/myfolder/myfile.csv"
    load_job = self.client.load_table_from_uri(
        uri, temp_table_ref, job_config=job_config
    )
    load_job.result()  # 等待加载完成

    # 2. 执行MERGE语句去重合并
    merge_query = f"""
        MERGE INTO `{self.table_ref["object"]}` AS target
        USING `{temp_table_ref.full_table_id}` AS source
        ON target.id = source.id
        WHEN NOT MATCHED THEN
            INSERT (id, name) VALUES (source.id, source.name)
        -- 如需更新已存在的记录,可添加:WHEN MATCHED THEN UPDATE SET name = source.name
    """

    query_job = self.client.query(merge_query)
    query_job.result()  # 等待合并完成

    # 可选:删除临时表(不需要保留则执行)
    self.client.delete_table(temp_table_ref)

方案2:追加后批量去重(适合增量小数据)

如果CSV是增量数据且重复比例低,可先追加数据到目标表,再通过BigQuery的窗口函数批量去重,操作简单且成本较低。

代码示例

def append_and_deduplicate():
    # 1. 先追加CSV数据到目标表
    job_config = bigquery.LoadJobConfig(
        schema=[
            bigquery.SchemaField("id", "INTEGER"),
            bigquery.SchemaField("name", "STRING"),
        ],
        source_format=bigquery.SourceFormat.CSV,
        skip_leading_rows=1,
        write_disposition=bigquery.WriteDisposition.WRITE_APPEND,
    )

    uri = "gs://mybucket/myfolder/myfile.csv"
    load_job = self.client.load_table_from_uri(
        uri, self.table_ref["object"], job_config=job_config
    )
    load_job.result()

    # 2. 执行去重查询,覆盖原表
    deduplicate_query = f"""
        CREATE OR REPLACE TABLE `{self.table_ref["object"]}` AS
        SELECT id, name
        FROM (
            SELECT 
                *,
                ROW_NUMBER() OVER (PARTITION BY id ORDER BY _PARTITIONTIME DESC) AS rn
            FROM `{self.table_ref["object"]}`
        )
        WHERE rn = 1
    """

    query_job = self.client.query(deduplicate_query)
    query_job.result()

注:如果是分区表,用_PARTITIONTIME保留最新记录;非分区表可替换为业务时间字段,或直接取任意一条重复记录。

方案3:外部表+物化视图(适合频繁更新的CSV)

如果CSV文件会定期更新,且需要实时查询去重后的数据,可创建外部表指向GCS的CSV,再通过物化视图自动维护去重数据。

步骤示例

  1. 创建外部表:
def create_external_table():
    external_config = bigquery.ExternalConfig("CSV")
    external_config.source_uris = ["gs://mybucket/myfolder/myfile.csv"]
    external_config.skip_leading_rows = 1
    external_config.schema = [
        bigquery.SchemaField("id", "INTEGER"),
        bigquery.SchemaField("name", "STRING"),
    ]

    table = bigquery.Table(self.client.dataset("your_dataset").table("external_csv"))
    table.external_data_configuration = external_config
    self.client.create_table(table)
  1. 创建去重物化视图(通过SQL执行):
CREATE MATERIALIZED VIEW `your_project.your_dataset.deduplicated_view`
AS
SELECT id, name
FROM (
    SELECT 
        *,
        ROW_NUMBER() OVER (PARTITION BY id ORDER BY _FILE_NAME DESC) AS rn
    FROM `your_project.your_dataset.external_csv`
)
WHERE rn = 1

物化视图会自动刷新,保证数据与源CSV同步,无需手动执行加载操作。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 23:01:08