每天收到一批新訂單,只要 INSERT 進資料庫就好了嗎?
第一天通常可以。第二天重跑同一個檔案,重複列就出現了; 上游更正金額後,舊資料沒有更新;遇到取消訂單,又不知道該刪哪一筆。
這篇拍拍君要用 DuckDB 的 MERGE INTO,把新增、更新、刪除整理成一條可以安全重跑的增量 ETL。
一. 先定義成功:不是「SQL 有跑完」 #
可靠的增量匯入,至少要回答五個問題:
- 同一個業務鍵在來源出現兩次時,哪一筆勝出?
- 舊事件晚到時,會不會覆蓋較新的狀態?
- 同一批重跑後,結果是否完全不變?
- 新增、更新、刪除各有幾列?
- 寫入失敗時,資料與稽核紀錄是否一起回滾?
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()
三層保護各有工作:
etl_runs.batch_id擋掉同批次重跑;updated_at擋掉舊事件與相同內容;- 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 與測試則保證重跑不會把資料愈修愈壞。
把這些責任分開後,批次失敗不再是半夜抽卡。 你會知道哪個契約壞了、影響幾列,以及重跑是否安全。