本日核心價值 (Core Focus): 以分塊 CSV 做 Extract,用 Pydantic 驗證 schema,依對照表讓 Codex 產出 Transform,再以冪等鍵 Load 進 PostgreSQL;沒有 checksum 筆數對上就不寫入,並提供 dry-run 預覽。
概念說明與實戰情境 (Overview)
遷移失敗通常不是「SQL 寫錯一句」,而是全檔一次灌入、欄位 silently 轉型、重跑產生重複列。大規模 ETL 應切成可重試的塊:Extract CSV → 用 Pydantic 驗證每一列 → Transform(由 mapping table 生成、程式碼進版控)→ Load 用 COPY 或 batched INSERT,並以業務冪等鍵(例如 legacy_id)去重。Codex 的工作是根據對照表寫 transform 函式,不是連上正式庫直接改資料。Load 前必須比對來源列數、通過驗證列數、將寫入列數;對不上就中止。dry-run 只印統計與樣本,不碰資料庫。
關鍵操作與範例 (Implementation & Example)
對帳用四個數字,缺一就停:
| 計數 | 意義 | 失敗時 |
|---|---|---|
extracted |
CSV 讀到的列 | 檔案空或編碼錯 |
valid |
Pydantic + transform 通過 | mapping 或髒資料 |
rejected |
未通過列 | 寫入 quarantine,不進正式表 |
loaded |
upsert 成功列 | 必須等於 valid |
建議目錄:
etl/
mapping.csv # source_col,target_col,rule
transform.py # Codex 產出,人 review 後合併
load_orders.py # extract / validate / load / checksum
data/orders.csv
對照表要經過與來源系統相同的欄位審查:哪個欄必填、哪個枚舉可對、哪個金額欄禁止四捨五入到整數。審查結果寫在表裡,而不是只存在聊天紀錄。Codex 讀表產出 transform_row 之後,用抽樣列做 golden test;抽樣應包含空值、未知 status、跨年日期。範例如下:
source_col,target_col,rule
OrderNo,legacy_id,strip; required
OrderDate,ordered_at,parse_date:%Y/%m/%d
Qty,qty,int; min=1
Amt,amount,decimal:2
Status,status,map:O=open;C=closed;X=cancelled
給 Codex 的 Prompt 要限制輸出範圍,避免它順便改 loader:
Read etl/mapping.csv and the first 20 rows of etl/data/orders.csv.
Write or update transform_row() in etl/transform.py only.
Keep function signature: def transform_row(raw: dict[str, str]) -> OrderRow.
Raise ValueError on invalid rows. Do not connect to the database.
本機產出後用 Day 16 的測試迴圈鎖住規則(非法日期、未知 status、空 OrderNo)。Loader 本體如下,含 dry-run 與 checksum:
from __future__ import annotations
import argparse
import csv
import hashlib
from decimal import Decimal
from pathlib import Path
from pydantic import BaseModel, Field, field_validator
from transform import transform_row
CHUNK = 5_000
class OrderRow(BaseModel):
legacy_id: str = Field(min_length=1, max_length=32)
ordered_at: str
qty: int = Field(ge=1)
amount: Decimal
status: str
@field_validator("status")
@classmethod
def status_ok(cls, v: str) -> str:
allowed = {"open", "closed", "cancelled"}
if v not in allowed:
raise ValueError(f"status {v!r} not in {allowed}")
return v
def file_checksum(path: Path) -> str:
h = hashlib.sha256()
with path.open("rb") as f:
for block in iter(lambda: f.read(1024 * 1024), b""):
h.update(block)
return h.hexdigest()
def extract_chunks(path: Path):
with path.open(newline="", encoding="utf-8-sig") as f:
reader = csv.DictReader(f)
batch: list[dict[str, str]] = []
for row in reader:
batch.append(row)
if len(batch) >= CHUNK:
yield batch
batch = []
if batch:
yield batch
def validate_transform(raw_rows: list[dict[str, str]]) -> tuple[list[OrderRow], int]:
ok: list[OrderRow] = []
rejected = 0
for raw in raw_rows:
try:
ok.append(transform_row(raw))
except Exception:
rejected += 1
return ok, rejected
def load_postgres(rows: list[OrderRow], dsn: str) -> int:
# Idempotent upsert on legacy_id. Prefer COPY into a staging table,
# then INSERT ... ON CONFLICT (legacy_id) DO UPDATE.
raise NotImplementedError
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("--input", default="etl/data/orders.csv")
parser.add_argument("--dsn", default="")
parser.add_argument("--dry-run", action="store_true")
args = parser.parse_args()
src = Path(args.input)
extracted = 0
valid = 0
rejected = 0
loaded = 0
sample: list[dict] = []
for chunk in extract_chunks(src):
extracted += len(chunk)
rows, bad = validate_transform(chunk)
rejected += bad
valid += len(rows)
if args.dry_run:
sample.extend(r.model_dump(mode="json") for r in rows[:3])
continue
loaded += load_postgres(rows, args.dsn)
checksum = file_checksum(src)
print(
{
"sha256": checksum,
"extracted": extracted,
"valid": valid,
"rejected": rejected,
"loaded": loaded,
"dry_run": args.dry_run,
"sample": sample[:5],
}
)
if extracted == 0:
raise SystemExit("no source rows")
if not args.dry_run and loaded != valid:
raise SystemExit(
f"checksum mismatch: loaded={loaded} valid={valid} extracted={extracted}"
)
if __name__ == "__main__":
main()
PostgreSQL 端用 staging + 冪等鍵,避免 batched INSERT 在重跑時重複:
CREATE TABLE IF NOT EXISTS orders (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
legacy_id text NOT NULL UNIQUE,
ordered_at date NOT NULL,
qty integer NOT NULL CHECK (qty >= 1),
amount numeric(12,2) NOT NULL,
status text NOT NULL,
source_sha256 text NOT NULL
);
-- Staging: COPY from a CSV produced after Pydantic validation, then:
INSERT INTO orders (legacy_id, ordered_at, qty, amount, status, source_sha256)
SELECT legacy_id, ordered_at, qty, amount, status, :sha256
FROM stg_orders
ON CONFLICT (legacy_id) DO UPDATE
SET ordered_at = EXCLUDED.ordered_at,
qty = EXCLUDED.qty,
amount = EXCLUDED.amount,
status = EXCLUDED.status,
source_sha256 = EXCLUDED.source_sha256;
被拒絕的列不要消失。每一塊在 rejected > 0 時寫出 etl/quarantine/chunk-{n}.csv(原始列加上錯誤訊息),方便業務補資料而不是叫工程師猜。COPY 適合已驗證、欄位已對齊的 staging 檔;來源仍可能含引號與換行時,先在 Python 寫出乾淨的 staging CSV 再 COPY,不要對原始上傳檔直接 COPY。batched INSERT 適用於數千列以下或必須逐列拿到主鍵的情境;十萬列以上優先 staging + 一組 INSERT ... SELECT。無論哪種,正式庫連線字串只存在環境變數或 CI secrets,不寫進 Codex Prompt。
操作順序固定:
python load_orders.py --dry-run:只印 extracted / valid / rejected 與樣本。COPY 進 staging,再 INSERT ... ON CONFLICT。extracted == valid + rejected,且 loaded == valid(upsert 列數以 xmax / n 回報為準,至少來源 valid 列都要能對上 legacy_id)。sha256 寫進行或 run log,之後同一檔重跑可證明輸入沒變。遷移不是一次作業而是可重跑的工作流:同一 legacy_id 重跑應收斂成同一列,而不是多一列。因此 load 路徑必須 upsert,dry-run 路徑必須能在沒有 DSN 的 CI 跑。把 transform_row 的單元測試(未知 status、空 OrderNo、非法日期)交給 Day 16 的迴圈去補,但 禁止 讓迴圈修改 load_orders.py 的 checksum 判斷——那是保險絲,不是業務規則。Windows 上大檔 CSV 注意 utf-8-sig 與換行;跑 COPY 仍建議在能連 PostgreSQL 的 Linux / 容器。Codex 可以幫忙寫 transform_row 與測試,但 load 開關、DSN、dry-run 預設值由人掌控。
百萬列不要一次切換。先用 1% 抽樣檔跑完整 Extract → dry-run → staging → upsert,對完四元組再放大。切換當日保留舊表只讀複本,新表以 source_sha256 標明這一批輸入;對帳失敗就停寫、不自動切連線。mapping 變更必須先改測試再改 transform,避免「Codex 改了規則、checksum 卻因列數碰巧相同而放行」。
dry-run 輸出要讓沒寫過程式的人看懂:用中文標出「來源列/通過列/拒絕列」,再附三筆樣本。rejected 超過 1% 就當失敗,不要「先灌再說」。staging 表每批用 run_id 區隔,失敗可整批清空 staging 而不動正式表。COPY 前後都數 COUNT(*),不要只看指令成功。若來源 CSV 每天增量,冪等鍵用 legacy_id;同一鍵後到的檔覆寫前一版,舊的 source_sha256 另留稽核表。讓 Codex 把 mapping 每一列變成一個測試案例,比「再請模型看一次 CSV」便宜也比較穩。Windows 匯出的 CSV 常帶 BOM 與 Big5,讀檔固定 utf-8-sig,必要時先轉碼再進 Pydantic。
注意事項與常見失敗 (Pitfalls)
INSERT:第一次 mapping 錯會寫進十萬筆垃圾。修法:--dry-run 為預設習慣;CI 只跑 dry-run 與單元測試。legacy_id + UNIQUE + ON CONFLICT。0、日期 2026/13/01 被吞掉。修法:Pydantic 嚴格驗證,失敗列計入 rejected,不要默默 drop 又不記數。COPY 成功不代表業務列數正確。修法:印 checksum 四元組,loaded != valid 就 SystemExit。transform.py;DDL 另檔、另審。CHUNK = 5000 分塊 yield。本日總結 (Takeaways)
COPY / batched upsert。transform_row,不把正式庫連線交給模型。legacy_id 這類業務鍵,而不是「檔案裡的第幾列」。extracted / valid / rejected 與樣本;沒有 checksum 對上就不要 load。明日預告 (Next)
明天把「結構化轉換」用在網頁內容:在合法授權前提下,用 BeautifulSoup / Puppeteer 抓頁,再用 LLM 依 schema 解析成 JSON,並把抓取與解析分開、快取 HTML。