DE Day 20 增量載入設計:水位標記與冪等
執行需求:CPU 可跑。今天是採集篇的收尾,也是管線篇的開門。我們把 Day 19 的分頁處理器升級為「增量載入」:每天只抓「上次沒抓過」的資料,並保證整個抓取流程是冪等(idempotent)的——重複執行會得到相同結果,不會造成資料重複或漏抓。這篇會建立兩個關鍵工程概念:水位標記(watermark)與冪等(idempotent),並用 DuckDB 1.4 把它們落到實處。版本基準仍是 2025 年 11 月:Python 3.13、DuckDB 1.4、httpx 0.28。
引言
Day 19 寫的分頁處理器可以一次抓完所有資料,但實務上我們每天要跑一次。一支每天跑的管線如果每次都「全部重抓」,會碰到三個問題:第一,浪費頻寬與時間,一個 10 萬筆的 API 全量重抓要 10 分鐘,增量只要 30 秒;第二,對 API 提供者不禮貌,特別是有每日配額的服務;第三,最重要的——若抓取中途失敗(例如網路斷線、API 變慢),重試時可能造成「同一天的資料被寫入兩次」,破壞資料完整性。
為了解決這些問題,資料工程發展出兩個互補的設計:水位標記(watermark)與冪等(idempotent)。水位標記是「我目前處理到哪裡」的紀錄,每次抓取只從這個位置往後取;冪等是「同一個操作重複執行結果相同」的特性,意思是即使同一筆資料被處理兩次,結果資料庫仍然只有一份。兩者結合起來,就能讓一支管線「每天跑、跑很久、隨時可重跑」。
今天的目標有四個:第一,定義水位標記與冪等的概念並用 SQL 表達;第二,設計「水位表」的 schema,存放在 DuckDB 的 raw._watermark;第三,把 Day 19 的分頁處理器升級為「增量版」,只在 fetched_at > watermark 時抓取;第四,用一段 SQL 驗證「重跑兩次結果一樣」。讀完之後你會知道為什麼 Day 1 寫的「水位標記、冪等」這兩個詞今天才展開——它們是資料管線從「能跑」走到「能上線」的關鍵設計。
水位標記與冪等的概念
先嚴謹定義這兩個概念。「水位標記(watermark)」是「資料流中目前已處理到的位置」。在時序資料裡,它通常是一個時間戳(例如「我已經處理到 2025-11-15 12:00 為止」);在編號資料裡,它通常是一個 ID(例如「我已經處理到訂單編號 12345 為止」);在檔案型資料裡,它通常是檔案的雜湊值或修改時間。每次抓取開始時,先讀水位標記;抓取結束時,更新水位標記。這個設計讓「從上次停的地方繼續」變得自然。
「冪等(idempotent)」是數學與電腦科學的標準術語,指「同一個操作執行任意次,結果都一樣」。在 HTTP 語意裡,GET、PUT、DELETE 是冪等的,POST 不是。在資料庫語意裡,「用相同 primary key 插入兩次」應該與「插入一次」結果相同,這通常靠「先刪除再插入」或「ON CONFLICT DO UPDATE」實作。在資料管線裡,冪等的意思是:同一段抓取腳本跑兩次,資料庫的內容應該一樣,不會多也不會少。
為什麼這兩個觀念一定要一起談?因為「水位標記」解決「要抓哪些」的問題,「冪等」解決「重抓會不會出包」的問題。如果只有水位標記但沒冪等,當水位標記因為某些原因被重設(譬如手動改回 0),重跑會把所有資料再抓一次並寫入資料庫,造成重複;如果只有冪等但沒水位標記,每天都要重抓所有資料,浪費資源。兩者合在一起才是完整的「增量載入」設計。
實務上常見的水位標記有三種:第一種是「時間戳水位」,例如 last_updated_at,適合「資料有更新時間欄位」的場景;第二種是「編號水位」,例如 last_id,適合「資料有單調遞增 ID」的場景;第三種是「檔案雜湊水位」,例如 last_md5,適合「每天會有新檔案、檔名有規律」的場景。今天的範例用第一種「時間戳水位」,因為它最容易理解也最常用。
完整實作:水位表 + 增量抓取 + 冪等驗證
今天的實作分五階段。第一階段建立水位表 schema;第二階段把 Day 19 的分頁處理器升級為「增量版」;第三階段設計冪等的寫入策略;第四階段實作驗證 SQL;第五階段用 DuckDB 跑一次完整的「重跑結果相同」測試。我們繼續沿用 pipelines/_mock_api.py(Day 19 的模擬 API)做示範。
第一步:建立水位表的 schema,把「每個資料源的水位」集中管理:
cd de-journey
mkdir -p warehouse
# 沿用 Day 19 的 _mock_api.py、_pager.py、_etiquette.py
uv pip install duckdb==1.4.1 polars==1.33.0
第二步:用一段 SQL 與 Python 程式建立水位表:
"""pipelines/_watermark.py:水位標記管理。"""
import duckdb
from datetime import datetime, timezone
SCHEMA = """
CREATE TABLE IF NOT EXISTS raw._watermark (
source_name VARCHAR PRIMARY KEY,
column_name VARCHAR NOT NULL,
last_value VARCHAR NOT NULL,
last_run_at TIMESTAMP NOT NULL,
n_rows INTEGER NOT NULL DEFAULT 0
)
"""
def init_watermark(con: duckdb.DuckDBPyConnection) -> None:
con.execute("CREATE SCHEMA IF NOT EXISTS raw")
con.execute(SCHEMA)
def get_watermark(con: duckdb.DuckDBPyConnection, source: str) -> tuple[str, str] | None:
row = con.execute(
"SELECT column_name, last_value FROM raw._watermark WHERE source_name = ?",
[source],
).fetchone()
return (row[0], row[1]) if row else None
def set_watermark(
con: duckdb.DuckDBPyConnection,
source: str,
column_name: str,
new_value: str,
n_rows: int,
) -> None:
con.execute("""
INSERT INTO raw._watermark (source_name, column_name, last_value, last_run_at, n_rows)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT (source_name) DO UPDATE SET
column_name = EXCLUDED.column_name,
last_value = EXCLUDED.last_value,
last_run_at = EXCLUDED.last_run_at,
n_rows = raw._watermark.n_rows + EXCLUDED.n_rows
""", [source, column_name, new_value, datetime.now(timezone.utc), n_rows])
if __name__ == "__main__":
con = duckdb.connect("warehouse/de-journey.duckdb")
init_watermark(con)
print("水位表已建立:raw._watermark")
set_watermark(con, "demo_api", "ts", "1730000300", 250)
print("示範水位寫入完成")
print(get_watermark(con, "demo_api"))
這段建立一個 raw._watermark 表,欄位含 source_name(資料源名稱)、column_name(水位欄位,例如 ts)、last_value(目前處理到的值)、last_run_at(上次執行時間)、n_rows(累計筆數)。用 ON CONFLICT DO UPDATE 保證冪等:同一個 source 第二次寫入會更新而非新增。set_watermark() 內的 n_rows + EXCLUDED.n_rows 是「累計本次新增」,讓水位表本身就是抓取歷史的摘要。
第三步:把 Day 19 的分頁處理器升級為「增量版」,抓取時讀水位、過濾資料、寫完更新水位:
"""pipelines/day20_incremental.py:增量版的分頁抓取器。"""
import asyncio
from datetime import datetime, timezone
from typing import Any
import duckdb
import httpx
from pipelines._etiquette import PoliteSession
from pipelines._pager import paginate_offset
from pipelines._watermark import get_watermark, set_watermark
UA = "HaoBot/1.0 (+https://blog.hao-code.com/bots)"
SOURCE = "demo_api"
URL = "http://127.0.0.1:8766/api/offset"
async def collect_incremental() -> int:
con = duckdb.connect("warehouse/de-journey.duckdb")
con.execute("""
CREATE TABLE IF NOT EXISTS raw.day20_items (
id INTEGER PRIMARY KEY,
name VARCHAR, ts BIGINT, fetched_at TIMESTAMP
)
""")
wm = get_watermark(con, SOURCE)
wm_ts = int(wm[1]) if wm else 0
print(f"水位標記:ts > {wm_ts}")
total = 0
async with httpx.AsyncClient() as client:
session = PoliteSession(UA, qps=3.0)
max_ts_seen = wm_ts
async for page in paginate_offset(client, URL, limit=50, session=session):
if not page.items:
break
fresh = [it for it in page.items if it["ts"] > wm_ts]
if not fresh:
continue
now = datetime.now(timezone.utc)
con.executemany(
"INSERT INTO raw.day20_items VALUES (?, ?, ?, ?) "
"ON CONFLICT (id) DO UPDATE SET "
"name = EXCLUDED.name, ts = EXCLUDED.ts, fetched_at = EXCLUDED.fetched_at",
[(it["id"], it["name"], it["ts"], now) for it in fresh],
)
total += len(fresh)
max_ts_seen = max(max_ts_seen, max(it["ts"] for it in fresh))
set_watermark(con, SOURCE, "ts", str(max_ts_seen), total)
return total
if __name__ == "__main__":
n = asyncio.run(collect_incremental())
print(f"本次新增:{n} 筆")
這段是今天的工程核心。重點有三:第一,進場時讀水位 wm_ts,只抓 ts > wm_ts 的資料;第二,寫入用 ON CONFLICT (id) DO UPDATE,保證冪等——即使同樣的 id 寫入兩次,結果還是同一筆;第三,寫完更新水位,把這次抓到的最大 ts 存回去。fresh = [it for it in page.items if it["ts"] > wm_ts] 這行做兩件事:先過濾掉舊資料,再用 con.executemany() 批次寫入。注意:我們在第一頁可能抓到的 ts 仍然可能小於水位(因為分頁中間沒有資料變動),這時 fresh 會是空 list,跳過寫入但仍繼續翻下一頁。
第四步:寫一段 SQL 驗證「重跑兩次結果相同」:
"""pipelines/day20_verify.py:驗證增量抓取是冪等的。"""
import duckdb
from pipelines.day20_incremental import collect_incremental
con = duckdb.connect("warehouse/de-journey.duckdb")
con.execute("""
CREATE TABLE IF NOT EXISTS raw.day20_verify (
run_id INTEGER, n_rows INTEGER, max_ts BIGINT
)
""")
def snapshot(run_id: int) -> tuple[int, int]:
n = con.execute("SELECT COUNT(*) FROM raw.day20_items").fetchone()[0]
max_ts = con.execute("SELECT COALESCE(MAX(ts), 0) FROM raw.day20_items").fetchone()[0]
con.execute("INSERT INTO raw.day20_verify VALUES (?, ?, ?)", [run_id, n, max_ts])
return n, max_ts
print("=== 第一次跑 ===")
import asyncio
asyncio.run(collect_incremental())
n1, t1 = snapshot(1)
print("=== 第二次跑(不應新增)===")
asyncio.run(collect_incremental())
n2, t2 = snapshot(2)
print(f"第一次:n={n1}, max_ts={t1}")
print(f"第二次:n={n2}, max_ts={t2}")
print(f"冪等驗證:{'通過' if (n1, t1) == (n2, t2) else '失敗'}")
這段示範「冪等驗證」的標準模式:跑第一次、抓快照、跑第二次、抓快照、比較兩次快照。第一次跑完應該抓完所有 250 筆,水位設到最大 ts;第二次跑時水位已經是最大 ts,所有資料都被過濾掉,n_rows 應該是 0(也就是沒新增)。但 raw.day20_items 內的筆數與最大 ts 應該完全不變,這就證明整個流程是冪等的。實際輸出取決於示範資料,第一次跑完 n=250,第二次 n 仍為 250,冪等驗證:通過。
第五步:把水位表與抓取狀態寫成可被外部查詢的「管線儀表板」SQL:
"""pipelines/day20_dashboard.py:抓取狀態儀表板。"""
import duckdb
con = duckdb.connect("warehouse/de-journey.duckdb")
print("=== 水位標記狀態 ===")
df = con.execute("""
SELECT source_name, column_name, last_value,
last_run_at, n_rows
FROM raw._watermark
ORDER BY source_name
""").fetchdf()
print(df.to_string(index=False))
print("=== 抓取歷史 ===")
df2 = con.execute("""
SELECT run_id, n_rows, max_ts, _loaded_at
FROM (
SELECT run_id, n_rows, max_ts FROM raw.day20_verify
UNION ALL
SELECT 0, COUNT(*), MAX(ts) FROM raw.day20_items
) t
ORDER BY run_id
""").fetchdf()
print(df2.to_string(index=False))
這段把水位表與抓取驗證表合併成一張「管線儀表板」SQL,輸出可以用 Jupyter notebook 或 Day 34 學的 Streamlit 儀表板顯示。每次跑抓取時,這張儀表板會自動更新;當 n_rows 突然下降、last_value 倒退、last_run_at 過期(譬如超過一天),就是警訊。
第六步:把水位設計的關鍵觀念寫成一份「增量載入檢查清單」,給接手的人當契約:
"""pipelines/_incremental_check.py:增量載入設計檢查清單(可執行)。"""
import json
CHECKLIST = {
"名稱": "DE Day 20 增量載入檢查清單",
"檢查": [
{"檢查": "水位表存在", "如何驗": "SELECT 1 FROM raw._watermark LIMIT 1"},
{"檢查": "資料源水位有寫入", "如何驗": "SELECT * FROM raw._watermark WHERE last_run_at > now() - interval '1 day'"},
{"檢查": "冪等鍵有 PRIMARY KEY 或 UNIQUE 約束", "如何驗": "DESCRIBE raw.day20_items(id 應為 INTEGER PRIMARY KEY)"},
{"檢查": "寫入用 ON CONFLICT DO UPDATE", "如何驗": "grep \"ON CONFLICT\" pipelines/day20_incremental.py"},
{"檢查": "重跑後筆數不變", "如何驗": "執行 day20_verify.py,第二次 n 應等於第一次 n"},
{"檢查": "水位前進(不倒退)", "如何驗": "SELECT last_value FROM raw._watermark 應單調遞增"},
],
}
if __name__ == "__main__":
print(json.dumps(CHECKLIST, ensure_ascii=False, indent=2))
這份「可執行清單」列出六個檢查點與對應的驗證指令。Day 38 的監控章節會把這份清單接到 CI 與排程中,每天自動跑一次並通知。執行這個檔案會印出 JSON 格式的檢查清單,方便對照設計是否符合預期。
第七步:用一段 SQL 把「水位表」與「驗證表」合併查詢,並把結果用 Polars 整理成可被儀表板吃的格式:
"""pipelines/day20_to_polars.py:把水位表轉 Polars DataFrame 給儀表板用。"""
import duckdb
import polars as pl
con = duckdb.connect("warehouse/de-journey.duckdb")
df = con.execute("""
SELECT
source_name,
column_name,
CAST(last_value AS BIGINT) AS last_value_bigint,
n_rows,
last_run_at
FROM raw._watermark
""").pl()
print(df)
print(f"資料源數:{df.height}")
print(f"總累計筆數:{df['n_rows'].sum()}")
# 找出超過 24 小時沒跑的資料源
stale = df.filter((pl.col("last_run_at") < pl.datetime(2025, 11, 15)) & (pl.col("last_run_at").is_not_null()))
print(f"可能過期的資料源:{stale['source_name'].to_list()}")
這段展示「把 DuckDB 查詢結果轉 Polars DataFrame」的標準做法。.pl() 是 DuckDB 內建的 Polars 整合介面,比 .df() 轉 pandas 更直接。Polars 的 filter() 用 expression 語法,比 SQL 更靈活。實際輸出會依水位表內容而定,示範情境應有 1 個資料源(demo_api),累計 250 筆。這份 Polars DataFrame 可以直接接到 Day 34 學的 Streamlit 儀表板,做成「資料源健康儀表板」。
常見錯誤與踩雷
第一個雷是「水位欄位選錯」。常見錯誤是用「資料建立時間」當水位,但實際上「資料更新時間」才是你要的。例如訂單資料的 created_at 是「訂單何時下單」,但「訂單狀態」可能後續更新(從「處理中」變「完成」),這時該用 updated_at。對應排查方向:在設計水位前先問「這份資料的變動性來自哪個欄位?」——如果是「狀態變動」就要找 update 欄位;如果是「新增資料」就用 create 欄位;如果有 soft delete 欄位(deleted_at),也要納入水位計算。
第二個雷是「時區沒統一」。水位用 UTC 存,但 API 回傳的是台灣時間,兩者差 8 小時。對應排查方向:所有時間一律存 UTC(TIMESTAMPTZ),顯示給使用者時再轉當地時區;水位計算時用 UTC 比較,避免「明明已經抓過,卻因為時差沒抓到」的問題。DuckDB 的 CAST(... AS TIMESTAMPTZ) 可以把 ISO 字串轉成 UTC 時間戳。
第三個雷是「水位表沒備份就改 schema」。水位表是管線的「記憶」,如果不小心刪了或改了欄位,整支管線會從頭跑。對應做法:水位表的 schema 變更要謹慎,最好寫成 migration;水位表本身也要定期 snapshot(譬如每週備份到另一個檔案),出問題時能還原。Day 38 的監控會把這條列入警報清單。
第四個雷是「冪等鍵選錯」。如果你用「資料的 hash 當 PK」,當資料稍微變動(譬如空白字元、大小寫),整筆會被當成新資料,造成重複與漏抓。對應做法:用「來源系統的業務鍵」當 PK(譬如訂單編號、會員 ID),不要用衍生欄位。如果來源沒有業務鍵,用「來源系統的主鍵欄位」加上一個「來源編號」組合成複合鍵。
第五個雷是「ON CONFLICT 沒生效」。PostgreSQL 與 DuckDB 的 ON CONFLICT 需要 PRIMARY KEY 或 UNIQUE 約束才會生效;如果沒設,ON CONFLICT 會被當成普通 INSERT,遇到重複會報錯。對應排查方向:先用 DESCRIBE 看表格的 PRIMARY KEY;如果沒有,先加 PRIMARY KEY 或 UNIQUE 約束;這部分在 Day 23 的維度建模會再展開。
效能與實務提醒
水位標記的儲存成本很低(每個資料源一筆),但更新頻率要看負載。如果每秒都要更新水位,可能會造成「水位表變成瓶頸」。對應做法:水位表寫入可以放寬到「批次更新」,譬如每抓 100 筆才更新一次;最後抓完再寫一次最終值。這個設計犧牲一點點「精確度」,換來「水位表不卡管線」。
另一個常見取捨是「過濾在客戶端 vs 伺服器端」。理想情況下,API 提供者會支援「過了某個時間只回新資料」的查詢(例如 ?since=2025-11-15),這時客戶端只需傳水位當參數,省下大量流量。但實務上很多 API 沒這功能,客戶端必須自己抓回來再過濾。本文的範例是後者,因為我們用的是 _mock_api.py,沒有 since 參數;真實工作時請優先用前者。
冪等的實作還有一個延伸:「schema 變動時怎麼辦」。如果資料源新增了一個欄位,舊的水位抓下來的資料會缺這個欄位。對應做法:在寫入前用 schema migration 工具(DuckDB 的 ALTER TABLE 或 dbt 的 schema 檔)先把表 schema 升級;水位表則記錄「上次成功的 schema 版本」,當 schema 變動時自動重置水位。這套設計在 Day 26 的 dbt 測試會更完整展開。
最後,水位標記與冪等的核心精神是「管線可以重跑」。一支好管線的標準是:今天跑失敗了,明天重跑可以從失敗的地方繼續;同一個抓取跑了兩次結果一樣。把這個精神寫進你的 README 與交接文件,是「管線活得比作者久」的關鍵。Day 22 的失敗處理會把這套概念延伸到「失敗通知」、「重跑工具」與「斷點續跑」,在那之前先把水位與冪等這層打穩。
小結
今天把 Day 19 的分頁處理器升級為「增量版」,建立了兩個資料工程的核心設計:水位標記(watermark)與冪等(idempotent)。我們用 DuckDB 實作了水位表 raw._watermark,用 ON CONFLICT DO UPDATE 保證冪等寫入,並用一段 SQL 驗證「重跑結果相同」。這套設計讓採集管線可以「每天跑、跑很久、隨時可重跑」,是整個系列從「能跑」走到「能上線」的關鍵工程。
把今天的關鍵詞整理進筆記本:水位標記(watermark)、冪等(idempotent)、ON CONFLICT DO UPDATE、PRIMARY KEY 約束、UTC 時間戳、增量載入、schema migration。這些詞在 Day 21(每日排程)、Day 22(失敗處理)、Day 30(端到端管線)會反覆出現。把 pipelines/_watermark.py、pipelines/day20_incremental.py、pipelines/day20_verify.py 都存起來,明天 Day 21 會把這套抓取邏輯掛到 APScheduler 的排程上。
結語
今天我們把「一次性抓取」變成「每天可重複的增量抓取」。但光有抓取邏輯不夠,還需要「每天準時執行」。明天 Day 21 我們會把這套邏輯掛到 APScheduler 3.11 上,建立一支「每天凌晨 3 點自動抓、跑完寫日誌、抓取失敗時留紀錄」的排程腳本,並討論排程的四個常見陷阱:時區、夏令時間、資源競爭、與單機限制。
明天,我們會把抓取管線從「手動跑」變成「自動跑」,這是資料工程從「個人工具」走到「團隊資產」的關鍵一步。
延伸資源
- DuckDB 1.4 ON CONFLICT 文件(2025):
https://duckdb.org/docs/sql/statements/insert.html,ON CONFLICT DO UPDATE的語法與限制。 - DuckDB TIMESTAMPTZ 文件(2025):
https://duckdb.org/docs/sql/data_types/timestamp,時區處理與 UTC 儲存的標準做法。 - Apache Airflow 增量載入模式(2025):
https://airflow.apache.org/docs/apache-airflow/stable/authoring-and-scheduling/timetable.html,增量管線在大型編排工具的標準設計。 - 政府資料開放平臺 API(2025):
https://data.gov.tw/api/v2/rest,真實世界的 CKAN 風格 API,依「政府資料開放授權條款第 1 版」可商用。 - 資料工程部落格「Designing Data-Intensive Applications」第十一章(Kleppmann, 2017),資料流系統中「水位標記」與「冪等」的學術基礎。
留言
張貼留言