iT邦幫忙

2026 iThome 鐵人賽

DAY 19
0
ChatGPT & Codex

ChatGPT + Codex 打造高效能 AI 開發工作流系列 第 19

Day 19: 大規模資料處理工作流:使用 Python/Codex 自動化 ETL 與 Data Migration

  • 分享至 

  • xImage
  •  

Day 19: 大規模資料處理工作流:使用 Python/Codex 自動化 ETL 與 Data Migration (ETL and Data Migration with Python Codex)

本日核心價值 (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。

操作順序固定:

  1. python load_orders.py --dry-run:只印 extracted / valid / rejected 與樣本。
  2. 人看 rejected 比例;過高就修 mapping / transform,不要 load。
  3. 正式跑時 COPY 進 staging,再 INSERT ... ON CONFLICT
  4. 比對 extracted == valid + rejected,且 loaded == valid(upsert 列數以 xmax / n 回報為準,至少來源 valid 列都要能對上 legacy_id)。
  5. 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)

  • 沒有 dry-run 就對正式庫 INSERT:第一次 mapping 錯會寫進十萬筆垃圾。修法:--dry-run 為預設習慣;CI 只跑 dry-run 與單元測試。
  • 用列號當唯一鍵:重跑或檔案合併會重複。修法:業務 legacy_id + UNIQUE + ON CONFLICT
  • 信任 CSV 型別:空字串變 0、日期 2026/13/01 被吞掉。修法:Pydantic 嚴格驗證,失敗列計入 rejected,不要默默 drop 又不記數。
  • Load 後不對筆數:COPY 成功不代表業務列數正確。修法:印 checksum 四元組,loaded != validSystemExit
  • 讓 Codex 直連正式庫改 schema:遷移腳本必須進 PR。修法:Codex 只改 transform.py;DDL 另檔、另審。
  • 單次讀完整檔進記憶體:數 GB CSV 會 OOM。修法:CHUNK = 5000 分塊 yield。

本日總結 (Takeaways)

  • ETL 以分塊進行:Extract → Pydantic 驗證 → Transform → COPY / batched upsert。
  • 對照表進版控,讓 Codex 只產出 transform_row,不把正式庫連線交給模型。
  • 冪等靠 legacy_id 這類業務鍵,而不是「檔案裡的第幾列」。
  • dry-run 先看 extracted / valid / rejected 與樣本;沒有 checksum 對上就不要 load。
  • 同一來源檔的 SHA-256 要留下,重跑時才能證明輸入與結果可對帳。

明日預告 (Next)

明天把「結構化轉換」用在網頁內容:在合法授權前提下,用 BeautifulSoup / Puppeteer 抓頁,再用 LLM 依 schema 解析成 JSON,並把抓取與解析分開、快取 HTML。


上一篇
Day 18: 實戰應用:Line Chatbot + AI Agent 打造企業級 ERP 查詢助手
下一篇
Day 20: 自動化 Web Scraping 工作流:結合 BeautifulSoup / Puppeteer 與 LLM 解析
系列文
ChatGPT + Codex 打造高效能 AI 開發工作流21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言