↓快轉到主要內容
  1. 教學文章/

DuckDB MERGE 實戰:增量匯入、去重與可重跑 ETL

·6 分鐘· loading · loading · ·
Python DuckDB SQL Merge ETL Data-Engineering Idempotency
每日拍拍
作者
每日拍拍
科學家 X 科技宅宅
目錄
Python 學習 - 本文屬於一個選集。
§ 136: 本文

featured

每天收到一批新訂單,只要 INSERT 進資料庫就好了嗎?

第一天通常可以。第二天重跑同一個檔案,重複列就出現了; 上游更正金額後,舊資料沒有更新;遇到取消訂單,又不知道該刪哪一筆。

這篇拍拍君要用 DuckDB 的 MERGE INTO,把新增、更新、刪除整理成一條可以安全重跑的增量 ETL。

一. 先定義成功:不是「SQL 有跑完」
#

可靠的增量匯入,至少要回答五個問題:

  1. 同一個業務鍵在來源出現兩次時,哪一筆勝出?
  2. 舊事件晚到時,會不會覆蓋較新的狀態?
  3. 同一批重跑後,結果是否完全不變?
  4. 新增、更新、刪除各有幾列?
  5. 寫入失敗時,資料與稽核紀錄是否一起回滾?

MERGE INTO 只能解一部分;來源去重、交易與測試,才會把流程補完整。

如果還不熟 DuckDB,先看 Python DuckDB 入門; 排名與移動視窗則在 Window Functions。 今天不重講 CSV/Parquet 分析,只處理增量資料的正確性。

二. 安裝與目標表
#

用 uv 建立環境:

uv init duckdb-merge-demo
cd duckdb-merge-demo
uv add "duckdb==1.5.5"
uv add --dev pytest

不用 uv 也可以:

python -m venv .venv
source .venv/bin/activate
python -m pip install "duckdb==1.5.5" pytest

這篇維護的是「每張訂單目前最新狀態」:

CREATE TABLE IF NOT EXISTS orders_current (
    order_id VARCHAR PRIMARY KEY,
    customer VARCHAR NOT NULL,
    amount DECIMAL(12, 2) NOT NULL,
    status VARCHAR NOT NULL,
    updated_at TIMESTAMP NOT NULL,
    source_batch VARCHAR NOT NULL
);

CREATE TABLE IF NOT EXISTS etl_runs (
    batch_id VARCHAR PRIMARY KEY,
    source_rows BIGINT NOT NULL,
    deduped_rows BIGINT NOT NULL,
    inserted_rows BIGINT NOT NULL,
    updated_rows BIGINT NOT NULL,
    deleted_rows BIGINT NOT NULL,
    finished_at TIMESTAMP NOT NULL DEFAULT current_timestamp
);

order_id 是業務鍵,一張訂單在目標表只有一列。 若需求是保留每次狀態變化,應改用 event table 或 SCD Type 2,別把兩種模型硬揉在一起。

batch_id 最好是上游提供的 immutable ID,或由檔案內容 hash 產生; 只用檔名,擋不住「同名檔案被換內容」。

三. 來源先標準化,不要直接 MERGE 原始檔
#

假設來源 CSV 長這樣:

order_id,customer,amount,status,updated_at,op,source_seq
PYPY-001,Maho,680.00,paid,2026-10-09 09:00:00,U,1
PYPY-002,Kurisu,420.00,new,2026-10-09 09:02:00,U,1
PYPY-001,Maho,720.00,paid,2026-10-09 09:05:00,U,2
PYPY-003,Daru,0.00,cancelled,2026-10-09 09:06:00,D,1

先建立 staging table,集中轉型與批次資訊:

CREATE OR REPLACE TEMP TABLE incoming_raw AS
SELECT
    trim(order_id) AS order_id,
    trim(customer) AS customer,
    amount::DECIMAL(12, 2) AS amount,
    lower(trim(status)) AS status,
    updated_at::TIMESTAMP AS updated_at,
    upper(trim(op)) AS op,
    source_seq::BIGINT AS source_seq,
    $batch_id::VARCHAR AS source_batch
FROM read_csv(
    $source_path,
    header = true,
    all_varchar = true
);

all_varchar = true 不是逃避型別,而是把轉型位置集中在 SQL。 壞資料會在 staging 階段明確失敗,不會被自動推論成意外的日期或數值。

參數也要和 SQL 分開傳:

con.execute(
    STAGE_SQL,
    {"batch_id": batch_id, "source_path": str(source_path)},
)

接著做品質閘門:

SELECT
    count(*) FILTER (WHERE order_id IS NULL OR order_id = '') AS bad_key,
    count(*) FILTER (WHERE updated_at IS NULL) AS bad_time,
    count(*) FILTER (WHERE op NOT IN ('U', 'D')) AS bad_op,
    count(*) FILTER (WHERE op = 'U' AND amount < 0) AS bad_amount
FROM incoming_raw;

不要用 WHERE 悄悄濾掉壞資料。 ETL 應該拒絕批次並報出計數,而不是產生一張看似成功但少資料的表。

四. 決定性去重:同秒也要有答案
#

同一批中 PYPY-001 出現兩次,不能期待檔案最後一列自然勝出。 SQL 的資料列沒有隱含順序,要把規則寫出來:

CREATE OR REPLACE TEMP TABLE merge_source AS
SELECT
    order_id,
    customer,
    amount,
    status,
    updated_at,
    op,
    source_batch
FROM incoming_raw
QUALIFY row_number() OVER (
    PARTITION BY order_id
    ORDER BY updated_at DESC, source_seq DESC
) = 1;

規則很明確:較新的 updated_at 勝出;時間相同時,較大的 source_seq 勝出。

若上游沒有 sequence,可用 immutable event ID。 真的什麼都沒有,至少用完整列 hash 做固定 tie-breaker; 但 hash 只能讓結果穩定,不能替你發明正確的業務順序。

五. MERGE:只讓較新的事件改資料
#

完整的新增、更新與顯式刪除如下:

MERGE INTO orders_current AS target
USING merge_source AS source
ON target.order_id = source.order_id

WHEN MATCHED
     AND source.op = 'D'
     AND source.updated_at >= target.updated_at
THEN DELETE

WHEN MATCHED
     AND source.op = 'U'
     AND source.updated_at > target.updated_at
THEN UPDATE SET
    customer = source.customer,
    amount = source.amount,
    status = source.status,
    updated_at = source.updated_at,
    source_batch = source.source_batch

WHEN NOT MATCHED
     AND source.op = 'U'
THEN INSERT (
    order_id, customer, amount, status, updated_at, source_batch
)
VALUES (
    source.order_id,
    source.customer,
    source.amount,
    source.status,
    source.updated_at,
    source.source_batch
)

RETURNING merge_action, order_id;

source.updated_at > target.updated_at 有兩個效果:

  • 較舊批次晚到時,不會覆蓋新狀態;
  • 同一份內容重播時,不會產生多餘 UPDATE。

也就是:

apply(batch, apply(batch, state)) == apply(batch, state)

這裡故意沒有用 WHEN NOT MATCHED BY SOURCE THEN DELETE。 增量批次只包含「有變更的列」,某張訂單沒出現,不代表它已刪除。 只有完整 snapshot 才能把缺席解讀成刪除;delta 請使用明確的 tombstone 事件。

RETURNING merge_action 會回傳 INSERT、UPDATE、DELETE,剛好可以做批次 QA。

六. 資料與稽核紀錄使用同一個 Transaction
#

Python 流程的核心如下:

from collections import Counter
from pathlib import Path

import duckdb


def run_batch(db_path: Path, source_path: Path, batch_id: str) -> dict[str, int]:
    con = duckdb.connect(db_path)
    try:
        con.execute("BEGIN TRANSACTION")

        already_done = con.execute(
            "SELECT count(*) FROM etl_runs WHERE batch_id = ?",
            [batch_id],
        ).fetchone()[0]
        if already_done:
            con.execute("ROLLBACK")
            return {"INSERT": 0, "UPDATE": 0, "DELETE": 0}

        con.execute(STAGE_SQL, {
            "batch_id": batch_id,
            "source_path": str(source_path),
        })
        source_rows = con.execute(
            "SELECT count(*) FROM incoming_raw"
        ).fetchone()[0]

        validate_source(con)
        con.execute(DEDUP_SQL)
        deduped_rows = con.execute(
            "SELECT count(*) FROM merge_source"
        ).fetchone()[0]

        changed = con.execute(MERGE_SQL).fetchall()
        counts = Counter(action for action, _order_id in changed)

        con.execute(
            """
            INSERT INTO etl_runs (
                batch_id, source_rows, deduped_rows,
                inserted_rows, updated_rows, deleted_rows
            ) VALUES (?, ?, ?, ?, ?, ?)
            """,
            [batch_id, source_rows, deduped_rows,
             counts["INSERT"], counts["UPDATE"], counts["DELETE"]],
        )
        con.execute("COMMIT")
        return dict(counts)
    except Exception:
        con.execute("ROLLBACK")
        raise
    finally:
        con.close()

三層保護各有工作:

  1. etl_runs.batch_id 擋掉同批次重跑;
  2. updated_at 擋掉舊事件與相同內容;
  3. transaction 讓目標資料與稽核列一起成功、一起失敗。

只靠 batch ID,遇到相同內容換 ID 就會失效; 只靠 row-level 條件,又無法回答這批何時跑過。兩層都要留。

七. Row Count QA:讓每次變更留下證據
#

每次執行至少保存五個數字:

指標 回答的問題
source_rows 原始批次讀到幾列?
deduped_rows 去重後剩幾個業務鍵?
inserted_rows 新業務鍵有幾個?
updated_rows 較新狀態更新幾個?
deleted_rows 明確刪除幾個?

再加幾個便宜的 invariant:

changed_rows = sum(counts.values())

assert deduped_rows <= source_rows
assert changed_rows <= deduped_rows
assert source_rows == 0 or deduped_rows > 0

實務上還應與近期 baseline 比較。 平常刪除率小於 1%,今天突然變成 70%,即使 SQL 沒報錯也值得隔離批次。

八. Watermark 是效能提示,不是唯一防線
#

資料變大後,可以只掃最近區間:

SELECT *
FROM read_parquet('data/orders/*.parquet')
WHERE updated_at >= $watermark - INTERVAL '1 day';

往回看一天,是為了容納遲到事件、時鐘偏差與補檔。 重疊區會重讀舊資料,沒關係:去重與 MERGE 時間條件會讓重跑安全。

不要只用嚴格的 updated_at > last_watermark。 同一 timestamp 的遲到事件可能被永久漏掉;更穩定的游標通常是時間加 sequence/event ID。

九. 冪等性測試:真的跑兩次
#

不要只看 SQL 覺得「應該可重跑」。把批次執行兩次,再比較目標內容:

def snapshot(con: duckdb.DuckDBPyConnection) -> list[tuple]:
    return con.execute("""
        SELECT order_id, customer, amount, status, updated_at, source_batch
        FROM orders_current
        ORDER BY order_id
    """).fetchall()


def test_same_batch_does_not_change_target(tmp_path: Path) -> None:
    db_path = tmp_path / "warehouse.duckdb"
    source = tmp_path / "batch.csv"
    source.write_text(SAMPLE_CSV, encoding="utf-8")

    initialize(db_path)
    run_batch(db_path, source, "batch-001")
    with duckdb.connect(db_path) as con:
        first = snapshot(con)

    second_counts = run_batch(db_path, source, "batch-001")
    with duckdb.connect(db_path) as con:
        second = snapshot(con)

    assert second == first
    assert second_counts == {"INSERT": 0, "UPDATE": 0, "DELETE": 0}

還要用新 batch ID 重播相同內容,證明不是只靠批次表短路:

def test_same_content_is_row_idempotent(tmp_path: Path) -> None:
    db_path = tmp_path / "warehouse.duckdb"
    source = tmp_path / "batch.csv"
    source.write_text(SAMPLE_CSV, encoding="utf-8")

    initialize(db_path)
    run_batch(db_path, source, "delivery-A")
    counts = run_batch(db_path, source, "delivery-B")

    assert counts == {}

最後再補三個測試:

  • 先載入新事件,再送舊事件:新狀態不可倒退;
  • 同一 timestamp 兩列:較大的 source_seq 必須勝出;
  • 故意讓 audit insert 失敗:orders_current 不可留下半套更新。

這三題測的是業務排序與 transaction 邊界,比「可以 connect」有價值得多。

十. 常見失敗模式
#

直接 MERGE 未去重來源
#

同一 target key 對到多個 source row,結果難以稽核。先壓成一列,並寫清楚 tie-breaker。

每次 MATCH 都 UPDATE
#

內容沒變也寫入,會製造虛假的更新計數。至少比較 updated_at,必要時再比較 payload hash。

把缺席解讀成 DELETE
#

delta 不是 snapshot。除非契約明確,否則只接受顯式 delete event。

MERGE 完才在另一個 transaction 寫 audit
#

程式若在中間中斷,資料已變但批次表沒記錄。兩者必須一起 commit。

十一. 上線前 Checklist
#

  • 業務鍵、事件時間與 delete 語意已定義
  • 來源先轉型、驗證,再決定性去重
  • 舊事件不會覆蓋新狀態
  • RETURNING merge_action 已納入 row-count QA
  • batch audit 與資料變更位於同一 transaction
  • 同 batch 與不同 batch ID 的重播都已測試
  • watermark 有 overlap window 或複合游標
  • 異常新增/刪除比例會告警或隔離

結語
#

MERGE INTO 不是一句神奇 upsert 咒語。

可靠的增量 ETL,是一組可以驗證的規則:staging 負責轉型與品質檢查, window function 負責去重,MERGE 負責變更,RETURNING 提供證據, transaction 與測試則保證重跑不會把資料愈修愈壞。

把這些責任分開後,批次失敗不再是半夜抽卡。 你會知道哪個契約壞了、影響幾列,以及重跑是否安全。

延伸閱讀
#

Python 學習 - 本文屬於一個選集。
§ 136: 本文

相關文章

DuckDB ASOF Join 實戰:時間序列對齊與最近事件查詢
·7 分鐘· loading · loading
Python DuckDB SQL ASOF Join Time-Series Temporal Data Data-Engineering
DuckDB Window Functions 實戰:排名、移動平均與分組分析
·8 分鐘· loading · loading
Python DuckDB SQL Window Functions Data-Analysis Analytics OLAP
Python fsspec 實戰:統一讀寫本機、S3、HTTP 與資料管線路徑
·7 分鐘· loading · loading
Python Fsspec Filesystem S3 Data-Engineering ETL
DuckDB 遠端 Parquet 實戰:S3/R2、httpfs、Secrets 與 Pushdown
·7 分鐘· loading · loading
Python DuckDB Parquet S3 Cloudflare R2 Httpfs Data-Engineering
Python PyArrow 實戰:Parquet、Schema 與跨工具資料交換
·8 分鐘· loading · loading
Python PyArrow Apache Arrow Parquet Data-Engineering ETL
Streamlit + DuckDB 實戰:本地資料查詢 Dashboard
·8 分鐘· loading · loading
Python Streamlit DuckDB SQL Dashboard Data-Analysis