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

BigQuery中MERGE语句为何偶发目标表重复记录?

问题:BigQuery MERGE语句导致目标表出现重复记录

我正在构建数据管道将数据导入BigQuery,将变更数据捕获(CDC)信息复制至BigQuery的临时表,随后使用MERGE语句将最新CDC记录通过插入、删除或更新操作同步至主表。按理解该方式应能避免目标表出现重复记录,但部分记录的最新CDC记录在目标表中被重复插入两次甚至更多。

用于同步的MERGE语句如下:

DECLARE p_run_id STRING;

DECLARE src_range STRUCT < date_min DATETIME,
date_max DATETIME >;

SET
  p_run_id = 'scheduled__2023-06-28T16:45:00+00:00';

--SET p_run_id='manual__2023-05-02T01:58:55.639096+00:00';
SET
  src_range =(
    SELECT
      STRUCT(
        MIN(date_min) AS date_min,
        MAX(date_max) AS date_max
      )
    FROM (
      SELECT
        MIN(end_time) AS date_min,
        MAX(end_time) AS date_max
      FROM
       dashboard.cdc_assignments
      where
        run_id = p_run_id
      UNION DISTINCT
      SELECT
        MIN(cdc_end_time_original) AS date_min,
        MAX(cdc_end_time_original) AS date_max
      FROM
        dashboard.cdc_assignments
      where
        run_id = p_run_id
  ));

MERGE dashboard.fact_assignments m USING (
  SELECT
    *
  EXCEPT
(row_num)
  FROM
    (
      SELECT
        *,
        ROW_NUMBER() OVER(
          PARTITION BY delta.assignment_seq_id
          ORDER BY
            delta.cdc_change_id DESC
        ) AS row_num
      FROM
        dashboard.cdc_assignments delta
      WHERE
        delta.run_id = p_run_id
    )
  WHERE
    row_num = 1
) d ON m.assignment_seq_id = d.assignment_seq_id
AND m.end_time >= src_range.date_min
and m.end_time <= src_range.date_max
 
WHEN NOT MATCHED
AND d.cdc_operation IN (2, 4) AND d.status != 'DEL' THEN
INSERT
  (
    assignment_seq_id,
    assignment_id,
    cdc_change_id,
    cdc_commit_date,
    employee_seq_id,
    ....
  )
VALUES
  (
    d.assignment_seq_id                   ,
    d.assignment_id                       ,
    d.cdc_change_id                       ,
    d.cdc_commit_date,
    d.employee_seq_id,
    ....
  )
  WHEN MATCHED
  AND ( d.cdc_operation = 1 OR d.status='DEL' )
  THEN 
    DELETE
  WHEN MATCHED
  AND d.cdc_operation = 4
  AND COALESCE(m.cdc_change_id, b'') < COALESCE(d.cdc_change_id, b'')
  THEN
    UPDATE
    SET
      assignment_seq_id = d.assignment_seq_id,
      assignment_id = d.assignment_id,
      cdc_change_id = d.cdc_change_id,
      cdc_commit_date = d.cdc_commit_date,
      employee_seq_id = d.employee_seq_id,
     .... 

问题现象

  • 重新加载特定文件的数据至cdc_assignments和fact_assignments表时,此前重复的记录未再出现。
  • 问题影响约5%的加载记录:每日加载约60万条分配记录,其中约2.5万条出现重复。
  • 部分分配记录有4条CDC记录(属正常情况),但最新CDC记录在目标表中存在2条重复;还有部分分配记录仅2条CDC记录,最终最新CDC记录在目标表中出现4条重复。
  • 出现问题的分配记录均来自同一加载至cdc_assignments表的文件,MERGE仅针对该文件内容执行,且这些记录未出现在其他文件中,不存在竞态条件。

疑问:MERGE语句是否并非防止目标表重复的正确方式?MERGE是否无法保证仅处理一次记录?


分析与解决方案

BigQuery的MERGE本身是原子性操作,能保证一次执行中不会产生重复,但你的问题出在匹配条件逻辑、源数据唯一性控制或任务幂等性上,以下是具体分析和解决办法:

1. 核心问题:MERGE匹配条件的逻辑漏洞

你当前的MERGE匹配条件包含m.end_time >= src_range.date_min AND m.end_time <= src_range.date_max,这会导致:

  • 如果目标表中已有同assignment_seq_id的记录,但end_time不在当前计算的src_range范围内,MERGE会判定为NOT MATCHED,触发插入新记录,最终导致同一assignment_seq_id对应多条不同end_time的重复记录。
  • 若任务因重试或误触发重复执行,而新插入的记录end_time恰好不在src_range内,会再次触发插入,加剧重复问题。

2. 源数据唯一性隐患

虽然你用ROW_NUMBER() OVER(PARTITION BY delta.assignment_seq_id ORDER BY delta.cdc_change_id DESC)取最新CDC记录,但如果存在同一assignment_seq_id下多条CDC记录的cdc_change_id值重复的情况,ROW_NUMBER()会生成多个row_num=1的记录,导致源数据中同一assignment_seq_id有多条记录,MERGE时每条都会触发插入操作,产生重复。

3. 解决办法

(1)修正MERGE匹配条件

如果assignment_seq_id是主表的业务唯一键,应移除匹配条件中的end_time范围限制,仅保留m.assignment_seq_id = d.assignment_seq_id,确保目标表中同一assignment_seq_id的记录能被正确匹配,避免不必要的插入:

MERGE dashboard.fact_assignments m USING (
  -- 源数据逻辑不变
) d ON m.assignment_seq_id = d.assignment_seq_id
-- 移除end_time范围条件

(2)确保源数据的唯一性

优化ROW_NUMBER()的排序逻辑,添加额外字段打破cdc_change_id重复的平局,保证每个assignment_seq_id仅返回一条最新记录:

SELECT
  *,
  ROW_NUMBER() OVER(
    PARTITION BY delta.assignment_seq_id
    ORDER BY delta.cdc_change_id DESC, delta.cdc_commit_date DESC
  ) AS row_num
FROM dashboard.cdc_assignments delta
WHERE delta.run_id = p_run_id

(3)添加幂等性控制

  • 在主表中对assignment_seq_id和cdc_change_id创建唯一约束,强制保证记录唯一性,即使MERGE尝试插入重复记录也会报错终止,避免脏数据。
  • 在源数据过滤中添加逻辑,跳过已同步过的CDC记录:
SELECT * EXCEPT(row_num)
FROM (
  -- 原ROW_NUMBER逻辑
)
WHERE row_num = 1
AND NOT EXISTS (
  SELECT 1 FROM dashboard.fact_assignments
  WHERE assignment_seq_id = d.assignment_seq_id
  AND cdc_change_id = d.cdc_change_id
)

(4)控制任务执行唯一性

在任务启动前检查当前p_run_id是否已处理过(比如维护一个任务执行日志表),若已处理则直接跳过,避免重复执行MERGE。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 10:37:49