跳到主要內容

DE Day 22 失敗處理:重試、警報與斷點續跑

DE Day 22 失敗處理:重試、警報與斷點續跑

執行需求:CPU 可跑。今天把 Day 21 的排程腳本升級為「具備自我修復能力」的管線。我們會設計「失敗分類 → 重試策略 → 警報通知 → 斷點續跑」四個層次的處理,並用 DuckDB 1.4 與 Python 標準函式庫實作一個最小可用版本。版本基準仍是 2025 年 11 月:Python 3.13、DuckDB 1.4、httpx 0.28、tenacity 9.0。讀完這篇,你會知道怎麼把一支「會跑」的工具變成「會救自己」的工具。

引言

Day 21 我們把抓取腳本掛上 APScheduler,每天自動跑。但一支「每天自動跑」的管線最怕的不是「沒跑」,而是「跑了但失敗」——而且「你不知道它失敗」。這是資料工程裡最貴的故障:抓不到資料沒人發現,儀表板上的數字停在三天前,主管問「為什麼報表沒更新」時你才慌忙查 log。一支成熟的管線必須能「自己發現失敗、自動重試、必要的時候主動通知」——這就是失敗處理(failure handling)的主題。

今天的設計分四層。第一層是「失敗分類」:把錯誤分成「暫時性」與「永久性」,對應不同的處理策略。第二層是「重試策略」:暫時性錯誤可以重試,但要聰明地重試(指數退避、加抖動)。第三層是「警報通知」:永久性錯誤或重試次數用盡時,主動通知負責人。第四層是「斷點續跑」:抓取中斷後,下次執行能從中斷點繼續,不會重抓或漏抓。這四層合起來,就是一套完整的「自我修復管線」。

今天的目標有四個:第一,建立「失敗分類」的對照表,把常見錯誤歸類為 transient(暫時性)或 permanent(永久性);第二,設計「重試 + 退避 + 抖動」的標準模組;第三,用 DuckDB 寫一個「重試佇列(retry queue)」資料表,把失敗的任務排隊等待下次重試;第四,用 DuckDB 寫一個「check-point」資料表,把抓取進度記錄下來,支援斷點續跑。我們沿用 Day 19 的 _mock_api.py 與 Day 20 的水位設計。

失敗分類:transient 與 permanent

「失敗分類」是失敗處理的第一步。常見的錯誤可以分成兩大類。第一類是 transient(暫時性):網路瞬斷、API 暫時 503、連線逾時——這些錯誤「稍等一下再試」通常會成功。第二類是 permanent(永久性):API 回 401(未授權)、404(資源不存在)、資料格式錯誤——這些錯誤「再試幾次」也不會成功,反而會浪費資源。設計重試策略的第一步就是「不要把兩種錯誤混在一起處理」。

實務上常見的 transient 錯誤有:網路 ConnectionError、TimeoutException、HTTP 408(Request Timeout)、429(Too Many Requests)、500(Internal Server Error)、502(Bad Gateway)、503(Service Unavailable)、504(Gateway Timeout)。常見的 permanent 錯誤有:HTTP 400(Bad Request,通常是程式 bug)、401(Unauthorized,API key 過期)、403(Forbidden,權限不足)、404(Not Found,URL 錯誤)、422(Unprocessable Entity,資料格式錯)。

另一個常被忽略的是「靜默錯誤」。API 回 200 但內容是錯誤訊息、JSON 解析後欄位缺失、CSV 編碼錯誤——這些不會丟例外,但資料是錯的。這類錯誤要靠 Day 14 學的「品質檢查」與 Day 32 學的「資料合約」抓。今天先處理「會丟例外」的錯誤類型,靜默錯誤留給後續章節。

失敗分類後,每類錯誤對應不同的處理策略。Transient 走「重試 + 退避」:第一次失敗等 1 秒重試、第二次 2 秒、第三次 4 秒,加 0 到 0.5 秒的隨機抖動避免雪崩。Permanent 走「通知 + 標記」:立刻把錯誤訊息寫進 log,並通知負責人;如果這個任務是批次的一部分,可以選擇「跳過繼續」或「整批停止」。這兩種策略的差別在於:transient 是「自動處理」、permanent 是「需要人類介入」。

完整實作:重試、佇列、斷點續跑

今天的實作分五階段。第一階段用 tenacity 9.0 設計「智慧重試」;第二階段設計「失敗分類」策略;第三階段用 DuckDB 建立「重試佇列」資料表;第四階段實作「斷點續跑」的 check-point 機制;第五階段把失敗訊息寫到 log 並設計通知。我們沿用 Day 19 的 _mock_api.py 做為抓取目標。

第一步:安裝 tenacity 與沿用前面的套件:

cd de-journey
uv pip install tenacity==9.0.0 httpx==0.28.1 duckdb==1.4.1
# 沿用 Day 17-21 的 _etiquette.py、_pager.py、_watermark.py、_scheduled_jobs.py

第二步:寫「失敗分類器」與「智慧重試」模組:

"""pipelines/_retry.py:失敗分類 + 智慧重試。"""
import asyncio
import random
from dataclasses import dataclass
from enum import Enum
from typing import Awaitable, Callable, TypeVar
import httpx
from tenacity import (
    AsyncRetrying, stop_after_attempt, wait_exponential_jitter,
    retry_if_exception_type,
)

T = TypeVar("T")

class ErrorKind(str, Enum):
    TRANSIENT = "transient"
    PERMANENT = "permanent"

@dataclass
class ClassifiedError(Exception):
    kind: ErrorKind
    original: Exception
    message: str

    def __str__(self) -> str:
        return f"[{self.kind.value}] {self.message}"

def classify_http_error(e: httpx.HTTPStatusError) -> ClassifiedError:
    status = e.response.status_code
    if status in (408, 429, 500, 502, 503, 504):
        return ClassifiedError(ErrorKind.TRANSIENT, e, f"HTTP {status} 暫時錯誤")
    return ClassifiedError(ErrorKind.PERMANENT, e, f"HTTP {status} 永久錯誤")

def classify_network_error(e: httpx.RequestError) -> ClassifiedError:
    return ClassifiedError(ErrorKind.TRANSIENT, e, f"網路錯誤: {type(e).__name__}")

async def retry_transient(
    fn: Callable[[], Awaitable[T]],
    max_attempts: int = 4,
) -> T:
    """只重試 transient 錯誤;permanent 直接拋出去。"""
    retryer = AsyncRetrying(
        stop=stop_after_attempt(max_attempts),
        wait=wait_exponential_jitter(initial=1, max=15),
        retry=retry_if_exception_type(ClassifiedError),
        reraise=True,
    )
    last_exc: ClassifiedError | None = None
    for attempt in retryer:
        with attempt:
            try:
                return await fn()
            except httpx.HTTPStatusError as e:
                ce = classify_http_error(e)
                last_exc = ce
                if ce.kind == ErrorKind.PERMANENT:
                    raise ce
                raise ce
            except httpx.RequestError as e:
                ce = classify_network_error(e)
                last_exc = ce
                raise ce
    if last_exc:
        raise last_exc
    raise RuntimeError("retry_transient 結束但無結果,邏輯錯誤")

if __name__ == "__main__":
    async def demo() -> None:
        async def flaky() -> int:
            if random.random() < 0.5:
                raise httpx.ConnectError("simulated")
            return 42
        result = await retry_transient(flaky)
        print(f"成功:{result}")
    asyncio.run(demo())

這段是今天的工程核心。ClassifiedError 是我們自訂的例外類型,把錯誤分成 TRANSIENT 與 PERMANENT。classify_http_error() 把 HTTP 狀態碼對應到錯誤類型;classify_network_error() 把網路錯誤(ConnectionError、TimeoutException 等)當作 transient。retry_transient() 是智慧重試的核心:它呼叫使用者傳入的函式,若丟出 ClassifiedError 就判斷要不要重試;transient 才重試,permanent 立刻拋出去。配合 tenacity 的 wait_exponential_jitter(initial=1, max=15),重試間隔會是「1-1.5, 2-2.5, 4-4.5, 8-8.5」秒,每次加 0.5 秒的隨機抖動。實際執行 demo() 約有 50% 機率成功、50% 機率在 4 次嘗試後失敗。

第三步:用 DuckDB 設計「重試佇列」資料表,把失敗的任務排隊等待下次重試:

"""pipelines/_retry_queue.py:失敗任務的持久化佇列。"""
from dataclasses import dataclass
from datetime import datetime, timezone
import json
import duckdb

SCHEMA = """
CREATE TABLE IF NOT EXISTS raw._retry_queue (
    task_id VARCHAR PRIMARY KEY,
    payload JSON,
    attempts INTEGER NOT NULL DEFAULT 0,
    last_error VARCHAR,
    last_error_kind VARCHAR,
    next_run_at TIMESTAMP NOT NULL,
    created_at TIMESTAMP NOT NULL
)
"""

@dataclass
class RetryTask:
    task_id: str
    payload: dict
    attempts: int
    last_error: str
    last_error_kind: str

def enqueue(con: duckdb.DuckDBPyConnection, task_id: str, payload: dict,
            error: str, kind: str, delay_sec: int = 60) -> None:
    con.execute(SCHEMA)
    con.execute("""
        INSERT INTO raw._retry_queue
            (task_id, payload, attempts, last_error, last_error_kind, next_run_at, created_at)
        VALUES (?, ?, 1, ?, ?, ?, ?)
        ON CONFLICT (task_id) DO UPDATE SET
            attempts = raw._retry_queue.attempts + 1,
            last_error = EXCLUDED.last_error,
            last_error_kind = EXCLUDED.last_error_kind,
            next_run_at = EXCLUDED.next_run_at
    """, [
        task_id,
        json.dumps(payload),
        error,
        kind,
        datetime.now(timezone.utc).replace(microsecond=0),
        datetime.now(timezone.utc),
    ])

def fetch_due(con: duckdb.DuckDBPyConnection, now: datetime | None = None) -> list[RetryTask]:
    if now is None:
        now = datetime.now(timezone.utc)
    rows = con.execute("""
        SELECT task_id, payload, attempts, last_error, last_error_kind
        FROM raw._retry_queue
        WHERE next_run_at <= ?
        ORDER BY next_run_at LIMIT 100
    """, [now]).fetchall()
    return [RetryTask(r[0], json.loads(r[1]), r[2], r[3], r[4]) for r in rows]

if __name__ == "__main__":
    con = duckdb.connect("warehouse/de-journey.duckdb")
    enqueue(con, "demo_task", {"url": "http://127.0.0.1:8766/api/offset"}, "test error", "transient")
    print("已加入重試佇列")
    for t in fetch_due(con):
        print(f"  待重試:{t.task_id}(attempts={t.attempts})")

這段建立 raw._retry_queue 資料表,欄位含 task_id(任務識別)、payload(任務內容,JSON 格式)、attempts(已重試次數)、last_error 與 last_error_kind(上次錯誤)、next_run_at(下次可重試時間)、created_at(任務建立時間)。用 ON CONFLICT DO UPDATE 保證冪等:同一個 task_id 重複 enqueue 會累加 attempts,不會新增多筆。fetch_due() 撈出「已到時間可重試」的任務,方便下次排程掃描。這個設計的好處是「即使程式崩潰,失敗的任務也留在 DuckDB,下次啟動時會自動撈出來重試」。

第四步:實作「斷點續跑」的 check-point 機制:

"""pipelines/_checkpoint.py:抓取進度 check-point,支援斷點續跑。"""
from dataclasses import dataclass
from datetime import datetime, timezone
import duckdb

SCHEMA = """
CREATE TABLE IF NOT EXISTS raw._checkpoint (
    job_name VARCHAR PRIMARY KEY,
    last_completed_offset BIGINT NOT NULL DEFAULT 0,
    last_completed_cursor VARCHAR,
    last_run_at TIMESTAMP NOT NULL,
    items_processed BIGINT NOT NULL DEFAULT 0,
    status VARCHAR NOT NULL DEFAULT 'idle'
)
"""

@dataclass
class Checkpoint:
    job_name: str
    last_completed_offset: int
    items_processed: int

def save_checkpoint(con: duckdb.DuckDBPyConnection, job_name: str,
                    offset: int, items: int, status: str = "running") -> None:
    con.execute(SCHEMA)
    con.execute("""
        INSERT INTO raw._checkpoint
            (job_name, last_completed_offset, items_processed, last_run_at, status)
        VALUES (?, ?, ?, ?, ?)
        ON CONFLICT (job_name) DO UPDATE SET
            last_completed_offset = EXCLUDED.last_completed_offset,
            items_processed = EXCLUDED.items_processed,
            last_run_at = EXCLUDED.last_run_at,
            status = EXCLUDED.status
    """, [job_name, offset, items, datetime.now(timezone.utc), status])

def load_checkpoint(con: duckdb.DuckDBPyConnection, job_name: str) -> Checkpoint | None:
    row = con.execute(
        "SELECT last_completed_offset, items_processed FROM raw._checkpoint WHERE job_name = ?",
        [job_name],
    ).fetchone()
    if row is None:
        return None
    return Checkpoint(job_name, row[0], row[1])

if __name__ == "__main__":
    con = duckdb.connect("warehouse/de-journey.duckdb")
    save_checkpoint(con, "demo_job", 50, 50)
    cp = load_checkpoint(con, "demo_job")
    print(f"check-point 載入:offset={cp.last_completed_offset}, items={cp.items_processed}")

這段建立 raw._checkpoint 表,記錄每個 job 的「目前處理到的 offset」與「已處理筆數」。save_checkpoint() 每抓一頁就呼叫一次,把當前進度寫進 DuckDB;load_checkpoint() 在抓取開始時讀進來,決定要從哪個 offset 繼續。這就是「斷點續跑」的標準做法:抓取中斷後,下次從 check-point 讀 offset,跳過已經抓過的部分。用 ON CONFLICT DO UPDATE 保證冪等——同一個 job 寫多次只會更新進度,不會新增記錄。注意這裡的 check-point 與 Day 20 的「水位標記」是兩個不同層次:水位是「資料層級」、check-point 是「任務層級」。

第五步:把失敗處理、排程、抓取整合成一支完整的「自我修復管線」:

"""pipelines/day22_resilient.py:具備失敗處理、警報、斷點續跑的抓取管線。"""
import asyncio
import json
from datetime import datetime, timezone
from pathlib import Path
import duckdb
import httpx
from pipelines._etiquette import PoliteSession
from pipelines._pager import paginate_offset
from pipelines._retry import retry_transient, ErrorKind
from pipelines._watermark import get_watermark, set_watermark
from pipelines._checkpoint import save_checkpoint, load_checkpoint

UA = "HaoBot/1.0 (+https://blog.hao-code.com/bots)"
SOURCE = "demo_api"
URL = "http://127.0.0.1:8766/api/offset"
JOB_NAME = "daily_offset_fetch"
ALERT_LOG = Path("logs/day22_alerts.jsonl")
ALERT_LOG.parent.mkdir(exist_ok=True)

def alert(level: str, message: str, **extra) -> None:
    rec = {"ts": datetime.now(timezone.utc).isoformat(), "level": level, "message": message, **extra}
    with ALERT_LOG.open("a", encoding="utf-8") as f:
        f.write(json.dumps(rec, ensure_ascii=False) + "\n")
    print(f"[{level.upper()}] {message}")

async def collect_with_retry() -> int:
    con = duckdb.connect("warehouse/de-journey.duckdb")
    con.execute("""
        CREATE TABLE IF NOT EXISTS raw.day22_items (
            id INTEGER PRIMARY KEY, name VARCHAR, ts BIGINT, fetched_at TIMESTAMP
        )
    """)
    cp = load_checkpoint(con, JOB_NAME)
    start_offset = cp.last_completed_offset if cp else 0
    alert("info", f"從 offset={start_offset} 開始", items_processed=cp.items_processed if cp else 0)

    async def fetch_page(client, offset):
        async with client.stream("GET", URL, params={"offset": offset, "limit": 50}) as r:
            r.raise_for_status()
            return r.json()

    total = 0
    try:
        async with httpx.AsyncClient() as client:
            session = PoliteSession(UA, qps=3.0)
            current = start_offset
            for page in paginate_offset(client, URL, limit=50, session=session, start_offset=start_offset):
                if not page.items:
                    break
                now = datetime.now(timezone.utc)
                con.executemany(
                    "INSERT INTO raw.day22_items VALUES (?, ?, ?, ?) "
                    "ON CONFLICT (id) DO UPDATE SET name=EXCLUDED.name, ts=EXCLUDED.ts",
                    [(it["id"], it["name"], it["ts"], now) for it in page.items],
                )
                total += len(page.items)
                current += len(page.items)
                save_checkpoint(con, JOB_NAME, current, (cp.items_processed if cp else 0) + total)
    except Exception as e:
        alert("error", f"抓取中斷:{type(e).__name__}: {e}")
        raise

    save_checkpoint(con, JOB_NAME, current, (cp.items_processed if cp else 0) + total, status="done")
    alert("info", f"抓取完成:{total} 筆", total=total)
    return total

if __name__ == "__main__":
    n = asyncio.run(collect_with_retry())
    print(f"本次新增:{n} 筆")

這段把今天所有概念整合起來。collect_with_retry() 開頭讀 check-point,決定從哪個 offset 繼續;每抓一頁就 save_checkpoint() 一次,把進度寫進 DuckDB;中途若丟例外,alert("error", ...) 把訊息寫進 logs/day22_alerts.jsonl,下次啟動時從 check-point 繼續抓。這個設計讓管線具備「中斷可繼續」的能力。

第六步:把失敗處理的決策邏輯寫成可被 lint 的規則,這對 Day 38 的監控很關鍵:

"""pipelines/_failure_decision.py:失敗處理決策樹(可執行)。"""
from dataclasses import dataclass
from enum import Enum

class Outcome(str, Enum):
    RETRY = "retry"
    SKIP = "skip"
    ALERT_AND_STOP = "alert_and_stop"

@dataclass
class Decision:
    outcome: Outcome
    delay_sec: int
    reason: str

def decide(kind: str, attempts: int, max_attempts: int = 3) -> Decision:
    if kind == "transient" and attempts < max_attempts:
        return Decision(Outcome.RETRY, delay_sec=2 ** attempts, reason="transient 未達上限")
    if kind == "transient":
        return Decision(Outcome.ALERT_AND_STOP, delay_sec=0, reason="transient 已達上限")
    if kind == "permanent":
        return Decision(Outcome.ALERT_AND_STOP, delay_sec=0, reason="permanent 永久錯誤")
    if kind == "silent":
        return Decision(Outcome.SKIP, delay_sec=0, reason="silent 靜默錯誤跳過")
    return Decision(Outcome.ALERT_AND_STOP, delay_sec=0, reason=f"未知錯誤類型 {kind}")

if __name__ == "__main__":
    for kind in ["transient", "permanent", "silent"]:
        for n in [0, 1, 2, 3]:
            d = decide(kind, n)
            print(f"kind={kind:10s} attempts={n} → {d.outcome.value:18s} ({d.reason})")

這個決策函式把失敗處理的三種策略(retry / skip / alert_and_stop)寫成可被 lint 的規則。transient 在 attempts 內會重試,超過就 alert_and_stop;permanent 立刻 alert_and_stop;silent 跳過並記錄。執行這個檔案會印出 12 行決策樹,方便對照邏輯是否符合預期。

第七步:寫一段簡單的「失敗演練」腳本,主動製造錯誤來測試我們的處理邏輯是否正確:

"""pipelines/day22_drill.py:失敗演練,確認管線在錯誤情境下仍可運作。"""
import asyncio
import httpx
from pipelines._retry import retry_transient, ErrorKind

async def simulate_transient_failure() -> None:
    """模擬一個 503 錯誤,確認會被分類為 transient 並重試。"""
    async def bad_call() -> str:
        r = httpx.Response(503)
        raise httpx.HTTPStatusError("503", request=httpx.Request("GET", "http://x"), response=r)
    try:
        await retry_transient(bad_call, max_attempts=2)
    except Exception as e:
        print(f"最終結果:{type(e).__name__}: {e}")

async def simulate_permanent_failure() -> None:
    """模擬一個 404 錯誤,確認會被分類為 permanent 且不重試。"""
    async def bad_call() -> str:
        r = httpx.Response(404)
        raise httpx.HTTPStatusError("404", request=httpx.Request("GET", "http://x"), response=r)
    try:
        await retry_transient(bad_call, max_attempts=5)
    except Exception as e:
        print(f"最終結果:{type(e).__name__}: {e}")

if __name__ == "__main__":
    asyncio.run(simulate_transient_failure())
    asyncio.run(simulate_permanent_failure())

這段是「失敗演練」——刻意製造錯誤,確認我們的處理邏輯會走到預期的路徑。simulate_transient_failure() 應該看到 retry 2 次後失敗(max_attempts=2),simulate_permanent_failure() 應該第一次就失敗(permanent 不重試)。把這種演練寫進 CI 與監控,可以確保「失敗處理」這段程式碼不會在某次改版後被悄悄弄壞。實務上,大型管線會把這種演練叫做「chaos testing」,每天隨機挑幾支 job 注入故障,確認 alert 會正確觸發。

另一個值得展開的議題是「分散式失敗處理」。今天我們設計的是「單機重試」,但實務上管線常跨多台機器或多個服務。例如一支管線可能從 RabbitMQ 拉任務、寫到 S3、通知 Slack,每一段都可能失敗。Day 27 的 Airflow 會處理這種分散式場景:每個 task 是獨立的,失敗由 Airflow 統一管理、重試、通知。今天的模組在這種架構下仍然適用——每個 task 內部的重試交給 tenacity,task 之間的重試交給 Airflow。

把失敗處理與通知的整合展開一下。今天我們用 alert() 把訊息寫進 JSON Lines log,正式部署時通常會接 Slack webhook:

  • 把 webhook URL 放進環境變數 SLACK_WEBHOOK_URL。
  • 把 alert() 改成「寫 log + POST 到 webhook」。
  • 把 alert 分級用 Slack 的 channel 區隔(critical → #data-alert、warning → #data-warning)。
  • 每天上午 9 點把昨天的 alert 摘要彙總寄 email,避免單則通知疲勞。

這套設計在 Day 37 的 GitHub Actions 與 Day 38 的監控章節會再展開。失敗處理不是單一技術,而是「觀察 → 分類 → 處理 → 通知 → 學習」的閉環;今天我們把前三個落地,明天開始的建模篇會把這個閉環延伸到資料層面。

常見錯誤與踩雷

第一個雷是「把所有錯誤都當 transient 重試」。很多新手會寫「try/except 包起來重試 5 次」,結果遇到 401(未授權)也重試 5 次,浪費 5 次配額後仍然失敗。對應排查方向:先做失敗分類,permanent 不重試;transient 才重試。實務上,常見的「不應該重試」錯誤有:400、401、403、404、422——這些通常是程式 bug 或設定錯誤,不是暫時問題。

第二個雷是「重試間隔太短」。如果第一次失敗立刻重試,等於對伺服器發出連環攻擊;如果一千個客戶端同時這麼做,就是 DoS 攻擊。對應排查方向:用指數退避(每次等更久),加抖動(避免雪崩)。tenacity 的 wait_exponential_jitter 是標準做法,本文示範用「initial=1, max=15」表示第一次等 1 秒、最大不超過 15 秒。

第三個雷是「重試次數無上限」。一支「理論上會無限重試」的腳本,會在伺服器真的當機時把整支管線卡住。對應做法:永遠設 max_attempts 上限(建議 3 到 5 次);達到上限就 alert_and_stop,交給人類判斷。

第四個雷是「失敗 log 沒寫就重試」。如果重試時前一次的錯誤訊息沒留下,事後追問題會找不到線索。對應做法:每次重試都把錯誤訊息、時間、attempts 寫進 log;log 格式用 JSON Lines,方便日後查詢(Day 32 會把 log 讀進 DuckDB 做分析)。

第五個雷是「斷點續跑沒對齊水位」。如果 check-point 寫的是 offset 50、水位寫的是 ts=1730000300,兩者不一致會造成資料錯亂。對應做法:check-point 與水位選一個當「主」、另一個當「副」。本文用 offset 當主、ts 當副(透過 fetched_at 欄位),實際工作時可以反過來,端看哪個指標更穩定。

效能與實務提醒

失敗處理對效能的影響主要在「重試等待」。一支抓 10 頁的腳本,若每頁重試 1 次等 1 秒,多花 10 秒;若全部重試 3 次等 7 秒,可能多花 70 秒。對應做法:根據「重要性」與「時效性」決定重試上限——離線報表可以寬鬆(即時性低、可以重試 5 次),即時儀表板要嚴格(即時性高、重試 2 次就 alert)。

另一個實務重點是「通知的去敏感」。一支每天 alert 10 次的管線,會讓負責人麻痺;一支一年只 alert 2 次的管線,反而會被認真看。對應做法:把 alert 分級——critical(資料完全沒抓到)、warning(部分資料缺漏)、info(執行成功但有警告);critical 才即時通知、warning 一天彙總一次、info 進 log 就好。本文用 alert("error", ...) 是 critical 等級,正式部署時可以接 Slack webhook(critical)或 email digest(warning)。

「斷點續跑」的設計還有一個延伸議題:「如果 check-point 寫了,但資料還沒寫進 DuckDB 怎麼辦?」這是經典的「分散式交易」問題。對應做法:把「寫 check-point」與「寫資料」放在同一個 DuckDB transaction 裡——要嘛都成功,要嘛都失敗。DuckDB 的 con.begin() 與 con.commit() 提供標準交易介面,但對單機抓取來說,「每頁一次 commit」已經夠安全。

最後,「自我修復管線」的核心精神是「讓程式具備處理失敗的能力,而不是讓程式逃避失敗」。失敗處理不是把錯誤吞掉,而是「把錯誤轉化為資訊」——什麼時候失敗、為什麼失敗、怎麼恢復——這些資訊最終會變成「管線活得比作者久」的關鍵。Day 38 的監控章節會把這套延伸到「自動告警、自動診斷、自動恢復」的完整閉環。

小結

今天把 Day 21 的排程腳本升級為「自我修復管線」。我們建立了四個關鍵設計:失敗分類(transient vs permanent)、智慧重試(tenacity 指數退避加抖動)、重試佇列(raw._retry_queue 持久化失敗任務)、斷點續跑(raw._checkpoint 記錄抓取進度)。這套設計讓管線可以「失敗自動重試、重試失敗自動通知、中途中斷下次繼續」。把今天的關鍵詞整理進筆記本:transient、permanent、ClassifiedError、指數退避、抖動、重試佇列、check-point、自我修復、alert 分級。

把 pipelines/_retry.py、pipelines/_retry_queue.py、pipelines/_checkpoint.py、pipelines/day22_resilient.py 都存起來。採集篇的 7 篇(Day 15-22)到這裡結束,明天 Day 23 我們進入「建模」篇,從資料倉儲的事實表與維度表開始,把今天抓回來的原始資料轉換成可分析的結構。

結語

今天我們把「會跑」變成「會救自己」。一支成熟的資料管線不是「不會失敗」,而是「失敗了知道怎麼救」。這個精神從 Day 22 開始延伸到整個系列——Day 30 的端到端管線、Day 32 的品質監控、Day 38 的資料合約,都會用到今天的失敗處理模組。明天開始,我們進入「建模」篇,把這些原始資料轉換成「事實表 + 維度表」的倉儲結構,讓分析師可以直接用 SQL 回答業務問題。

明天,我們會用 DuckDB 把「抓回來的訂單資料」轉成維度模型,並用 dbt 工具管理這個轉換流程。

延伸資源

  • tenacity 9.0 文件(2025):https://tenacity.readthedocs.io/en/latest/index.html,wait_exponential_jitter、retry_if_exception_type 的標準用法。
  • Microsoft Azure Architecture「Retry pattern」(2025):https://learn.microsoft.com/en-us/azure/architecture/patterns/retry,失敗分類與重試策略的權威參考。
  • DuckDB 1.4 Transaction 文件(2025):https://duckdb.org/docs/sql/transactions.html,BEGIN、COMMIT、ROLLBACK 的標準用法。
  • 政府資料開放平臺 API(2025):https://data.gov.tw/api/v2/rest,真實世界的 CKAN 風格 API,依「政府資料開放授權條款第 1 版」可商用。
  • 《Designing Data-Intensive Applications》第十一章(Kleppmann, 2017),資料流系統中 failure handling 與 check-point 的學術基礎。

留言

這個網誌中的熱門文章

Day 2 變數與資料型別

Day 2 變數與資料型別 引言 寫程式的過程中,變數與資料型別是處理資料的基礎。變數是存放資料的容器,資料型別則決定這筆資料有哪些特性、可以進行哪些操作。學會定義變數、認識各種資料型別,是學好 Python 的關鍵一步。 這篇文章會帶你了解 Python 中變數的觀念、如何定義變數,以及常見的資料型別,包括整數、浮點數、字串、布林值,還有串列、元組、字典與集合等容器型別。我們也會介紹變數的命名規則與撰寫風格建議,以及如何用 type() 檢查資料型別。 什麼是變數?如何在 Python 中定義變數 變數是在程式執行時用來存放資料的名稱。透過定義變數,我們可以給一筆資料一個名字,並在程式的其他地方用這個名字取用該筆資料。在 Python 中,變數不需要事先宣告型別,因為 Python 是動態型別語言,變數的型別由指定給它的值決定。 定義變數的基本語法 在 Python 中定義變數非常簡單,只要用賦值符號 = 把值指定給變數即可。例如: x = 5 # 定義變數 x,並把整數 5 賦值給它 name = "Alice" # 定義變數 name,並把字串 "Alice" 賦值給它 在這裡,x 是一個變數,被賦予整數 5;name 是另一個變數,被賦予字串 "Alice"。 變數的更新與覆寫 變數的值可以修改,也就是說,我們可以在程式的不同地方給同一個變數新的值。例如: x = 10 # x 最初被賦予 10 x = 15 # x 的值現在被更新為 15 這樣就能依照需求,在程式執行過程中靈活調整變數的值。 Python 的動態型別系統 Python 和某些靜態型別語言不同,定義變數時不需要宣告型別。賦值時,Python 會根據值自動判斷變數的型別。例如: x = 5 # x 是整數 x = 3.14 # x 變成浮點數 x = "Hi" # x 變成字串 同一個變數在程式執行過程中可以存放不同型別的值,這是 Python 的彈性之一。 常見資料型別 在 Python 中,資料型別決定我們可以對變數進行哪些操作...

Day 1 Python 簡介與環境設定

Day 1 Python 簡介與環境設定 引言 在現在的科技環境裡,程式設計已經是一項重要技能。無論你是對資料科學有興趣、想成為開發者,或是想踏入人工智慧(AI)領域,學會寫程式都能明顯提升你的競爭力。在眾多程式語言中,Python 因為語法簡單、功能強大、應用範圍廣泛,成為許多人進入程式世界的第一選擇。這篇文章會帶你認識 Python 的背景與優勢,並一步步教你在不同系統上安裝與設定 Python 開發環境,最後寫出第一支 Python 程式。 為什麼選擇 Python? Python 是一種高階程式語言,由 Guido van Rossum 在 1991 年發布。Python 的設計哲學強調程式碼的可讀性,並用縮排來定義程式區塊,這點和許多使用大括號的語言不同。簡潔的語法讓它成為初學者的理想選擇;就算是經驗豐富的開發者,也能用它完成複雜的專案。 Python 的優勢如下: 簡單易學 :Python 的語法清楚、結構簡潔,初學者很快就能上手。和其他語言相比,學習曲線相對平緩,不需要先弄懂一堆複雜觀念,就能開始寫程式。 應用範圍廣泛 :從資料科學、網頁開發、人工智慧、機器學習、自動化測試到網路爬蟲,Python 都有大量開源函式庫與工具支援,而且在這些領域都扮演關鍵角色。 豐富的函式庫與框架 :Python 的函式庫生態系非常龐大。做資料分析有 NumPy、Pandas;開發網站有 Django、Flask;做深度學習有 TensorFlow、PyTorch。各種需求幾乎都能找到對應的套件,讓開發更有效率。 跨平台支援 :Python 支援 Windows、macOS、Linux 等作業系統,程式通常不需要太多修改就能跨平台執行,讓開發與部署更有彈性。 活躍的社群 :Python 擁有龐大的開發者社群。學習或開發上遇到問題,幾乎都能在社群與論壇(例如 Stack Overflow)找到答案,對初學者來說是很強的後盾,也能減少卡關時的挫折感。 Python 的應用領域 Python 的流行與強大功能,讓許多領域都開始大量使用它。以下是幾個常見的應用方向: 資料科學 :隨著大數據與人工智慧興起,資料科學大量使用 Python。NumPy、Pandas 與 Matplotlib 等工具能處理和分析龐...

Python 從入門到 PyTorch 深度學習:開啟 AI 世界的大門

Python 從入門到 PyTorch 深度學習:開啟 AI 世界的大門 隨著人工智慧(AI)與深度學習(Deep Learning)快速發展,越來越多人對這些技術產生興趣。不論你是想踏入 AI 領域的初學者,還是已經有程式基礎的開發者,學好 Python 與深度學習框架(例如 PyTorch),都能為你打開更多可能。 為什麼選擇 Python? Python 已經是資料科學與人工智慧領域的首選語言。它的語法簡潔、容易上手,而且擁有龐大的生態系與大量開源函式庫。無論是資料處理、資料視覺化,還是建立機器學習與深度學習模型,Python 都能勝任。對想進入 AI 或資料科學領域的人來說,它幾乎是必備工具。 PyTorch 是什麼? PyTorch 是由 Meta(原 Facebook)AI 研究團隊開發的開源深度學習框架,以易用、靈活和動態計算圖著稱,是許多 AI 研究人員與開發者的首選。相較於其他框架,PyTorch 的寫法更貼近原生 Python,對初學者相對友善。無論是簡單的實驗,還是複雜的深度學習模型,PyTorch 都能提供強大的支援。 這個系列能帶給你什麼? 這個系列會從 Python 的基礎開始,帶你一步一步學習,最後能自己用 PyTorch 建立深度學習模型。即使你完全沒有寫過程式,也能跟著文章的節奏累積技能,理解 AI 與深度學習的核心觀念。 本系列涵蓋的主題 Python 基礎:從變數、條件判斷到函式與模組。 資料處理工具:用 NumPy 與 Pandas 有效率地操作資料。 資料視覺化:用 Matplotlib 與 Seaborn 把資料畫成圖表。 深度學習的數學基礎:線性代數、微積分與機率。 PyTorch 入門:理解張量、模型建構與 GPU 加速。 基礎深度學習模型:CNN 與 RNN 的實作應用。 深度學習專案實戰:從資料前處理到模型部署的端到端流程。 誰適合這個系列? 程式初學者 :如果你對 AI 充滿好奇,卻還沒寫過程式,系列的第一部分會帶你快速上手 Python,並幫助你理解深度學習的基本觀念。 資料科學愛好者 :如果你已經熟悉一些資料處理方法,進階部分會教你如何用 PyTorch 建構深度學習模型。 開發者與研究人員 :想更深入了...