DE Day 19 API 資料蒐集:分頁、重試與增量
執行需求:CPU 可跑。今天進入採集篇的第四篇,把焦點從「抓整包」轉到「分頁抓 API」。真實世界的 API 很少一次回傳所有資料,常見的做法是用分頁(pagination)控制單次回應大小。這篇會把 Day 17 的禮貌套件與 Day 18 的 DuckDB 落地擴充成「分頁處理器」,用一個本地模擬的「指標查詢 API」示範兩種主流分頁——offset-based 與 cursor-based——並示範水位標記(watermark)的暖身動作:增量抓取。版本基準仍是 2025 年 11 月:Python 3.13、httpx 0.28、DuckDB 1.4。
引言
Day 18 抓政府開放資料時,我們一次下載一個檔案(CSV 或 JSON)。但很多 API 不是這樣設計的——他們一次只回 100 筆、1,000 筆,要拿全部資料必須「翻下一頁」。這在資料工程是個獨立的小學問,叫做分頁(pagination)。分頁的設計直接影響抓取速度、增量載入、與水位標記,今天就把它講清楚。
今天的目標有四個:第一,了解兩種主流分頁:offset-based(用「從第 N 筆開始、取 M 筆」)與 cursor-based(用「上次最後一筆的某個值」當起點);第二,寫一支通用的分頁處理器,能適應兩種風格;第三,用「指數退避 + jitter」處理 API 限流(rate limit)與暫時錯誤;第四,用水位標記的概念做增量抓取,避免每天重抓整個資料集。這四件事合起來,就是一支「可以每天跑、能跑很久、不會把對方打掛」的 API 採集管線。
先說明 API 採集的倫理。Day 17 講的 robots.txt、節流、User-Agent 仍然適用,但 API 還多了一條:「遵守服務條款與 API 配額」。多數 API 提供者會在「API 文件」或「服務條款」中明列配額(例如每分鐘 60 次、每天 10000 次)。這個配額是法律契約的一部分,超過就可能違反條款。實務上我們會用禮貌套件的節流來保守估計,並在程式裡寫明對方給的配額是多少;萬一配額改變而我們沒注意,後果由我們承擔。另一個倫理議題是「資料授權」,Day 18 已經處理過政府資料,這裡要追加的是「商業 API 通常有專屬授權條款」,不能直接拿來公開散布,要看清楚再用。
分頁的兩種主流設計
offset-based pagination 是最直覺的寫法:API 用 ?offset=0&limit=100 兩個參數告訴你「從第幾筆開始、取幾筆」。伺服器端簡單,呼叫端也簡單,但有三個問題:第一,越後面的頁面成本越高(伺服器通常要從頭數到 offset 那一筆),所以「拿到後段資料」會比較慢;第二,若資料在抓取過程中新增或刪除,offset 會錯亂,可能漏抓或重抓;第三,總頁數在資料成長時會變動,沒辦法事先知道。
cursor-based pagination 是另一種風格:API 用 ?cursor=abc123&limit=100 或 ?after_id=12345&limit=100,告訴你「從這個 ID 之後開始、取幾筆」。這個設計的好處是:每頁成本都一樣(伺服器只要找 cursor 對應的索引位置,往後取 limit 筆);新增或刪除不會錯亂(cursor 是絕對位置);對無限滾動的資料流(如社群媒體貼文、感測器讀數)特別適合。代價是:呼叫端不能「跳頁」,必須一頁一頁往後走;對應用程式來說,offset 比較直覺,但對「資料管線」來說,cursor 是更好的設計。
實務上還會碰到第三種:link-based pagination。回應裡直接給你「下一頁 URL」(例如 GitHub API 的 Link: <https://...>; rel="next" header)。這本質上是 cursor 的一種實作,只是把 cursor 藏在 URL 裡。還有第四種:time-based pagination,用「起始時間 + 結束時間」把資料切成區段,例如 ?start=2025-11-01&end=2025-11-02。這種在「時序資料」特別常見,例如感測器讀數、log 檔。下表把四種做個整理:
四種分頁風格對照:
- offset-based:
?offset=0&limit=100,直覺但慢,越後面越慢。 - cursor-based:
?cursor=abc123&limit=100,穩定,無限資料流適用。 - link-based:回應 header 給下一頁 URL,省得自己組。
- time-based:
?start=2025-11-01&end=2025-11-02,時序資料常用。
今天的程式碼會實作 offset 與 cursor 兩種,第三種只要包裝成 cursor 一樣能用,第四種留給 Day 20 的水位標記篇章。
完整實作:分頁處理器 + 限流 + 增量
今天的實作目標是建立一個「通用分頁處理器 pipelines/_pager.py」,能處理 offset-based 與 cursor-based 兩種 API,並整合 Day 17 的禮貌節流、Day 18 的 DuckDB 落地。為了不依賴特定 API(避免哪天服務改版範例就壞),我們用 FastAPI 寫一個迷你伺服器模擬兩種分頁風格,再用分頁處理器去打它。
第一步:安裝套件並沿用工作目錄:
cd de-journey
uv pip install httpx==0.28.1 fastapi==0.115.0 uvicorn==0.32.0
# 沿用 _etiquette.py 與 _robots.py(Day 17)
第二步:寫一個迷你 API 伺服器,模擬 offset-based 與 cursor-based 兩種分頁風格:
"""pipelines/_mock_api.py:模擬兩種分頁風格,給分頁處理器範例用。"""
import time
from fastapi import FastAPI, HTTPException, Query
from fastapi.responses import JSONResponse
app = FastAPI()
DATA = [{"id": i, "name": f"item-{i:04d}", "ts": 1730000000 + i * 60} for i in range(1, 251)]
@app.get("/api/offset")
async def offset_endpoint(offset: int = Query(0, ge=0), limit: int = Query(50, le=200)):
if offset >= len(DATA):
return {"items": [], "next_offset": None, "total": len(DATA)}
end = min(offset + limit, len(DATA))
next_offset = end if end < len(DATA) else None
return {"items": DATA[offset:end], "next_offset": next_offset, "total": len(DATA)}
@app.get("/api/cursor")
async def cursor_endpoint(cursor: int = Query(0, ge=0), limit: int = Query(50, le=200)):
if cursor >= len(DATA):
return {"items": [], "next_cursor": None}
end = min(cursor + limit, len(DATA))
next_cursor = end if end < len(DATA) else None
return {"items": DATA[cursor:end], "next_cursor": next_cursor}
@app.get("/api/slow")
async def slow_endpoint():
time.sleep(0.05)
return {"items": DATA[:5]}
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="127.0.0.1", port=8766)
這段建立兩個端點:/api/offset 接受 offset 與 limit,回傳資料切片與下一個 offset;/api/cursor 接受 cursor 與 limit,回傳下一個 cursor。用 uv run python pipelines/_mock_api.py 啟動它,瀏覽器打 http://127.0.0.1:8766/api/offset?offset=0&limit=10 可以看到第一頁資料。
第三步:寫通用分頁處理器 _pager.py:
"""pipelines/_pager.py:通用分頁處理器,支援 offset 與 cursor。"""
import asyncio
from dataclasses import dataclass
from typing import AsyncIterator, Callable
import httpx
from pipelines._etiquette import PoliteSession, AuditRecord
@dataclass
class PageRequest:
url: str
params: dict
@dataclass
class Page:
items: list
next_request: PageRequest | None
async def paginate_offset(
client: httpx.AsyncClient,
base_url: str,
limit: int,
session: PoliteSession,
start_offset: int = 0,
) -> AsyncIterator[Page]:
offset = start_offset
while True:
req = PageRequest(url=base_url, params={"offset": offset, "limit": limit})
r = await session.get(client, req.url)
r.raise_for_status()
body = r.json()
items = body.get("items", [])
yield Page(items=items, next_request=None)
next_offset = body.get("next_offset")
if not items or next_offset is None:
break
offset = next_offset
async def paginate_cursor(
client: httpx.AsyncClient,
base_url: str,
cursor_param: str,
limit: int,
session: PoliteSession,
start_cursor: int = 0,
) -> AsyncIterator[Page]:
cursor = start_cursor
while True:
params = {"cursor": cursor, "limit": limit}
r = await session.get(client, base_url)
r.raise_for_status()
body = r.json()
items = body.get("items", [])
yield Page(items=items, next_request=None)
next_cursor = body.get("next_cursor")
if not items or next_cursor is None:
break
cursor = next_cursor
if __name__ == "__main__":
async def demo() -> None:
async with httpx.AsyncClient() as client:
session = PoliteSession("HaoBot/1.0 (+https://blog.hao-code.com/bots)", qps=5.0)
count = 0
async for page in paginate_offset(client, "http://127.0.0.1:8766/api/offset", limit=50, session=session):
count += len(page.items)
print(f"offset 風格:抓到 {count} 筆")
asyncio.run(demo())
這段設計的重點在「用 async generator 把分頁邏輯與使用邏輯分離」。paginate_offset 與 paginate_cursor 兩個函式各自處理自己的分頁規則,但都回傳 AsyncIterator[Page],讓使用端用 async for page in paginate_xxx(...) 統一處理。每次 yield 一頁,呼叫端可以選擇「全部收集」、「逐頁寫 DuckDB」、「逐頁通知 Slack」等不同策略。禮貌節流與稽核則透過注入 PoliteSession 與 AuditRecord 達成,符合 Day 17 的設計精神。實際抓取 _mock_api.py 會跑完所有頁,總筆數應該是 250。
第四步:用分頁處理器抓資料並寫進 DuckDB,示範「批次寫入」與「增量標記」:
"""pipelines/day19_collect.py:用分頁處理器抓資料並落進 DuckDB。"""
import asyncio
from datetime import datetime
import duckdb
import httpx
from pipelines._etiquette import PoliteSession
from pipelines._pager import paginate_offset
UA = "HaoBot/1.0 (+https://blog.hao-code.com/bots)"
async def collect() -> int:
con = duckdb.connect("warehouse/de-journey.duckdb")
con.execute("CREATE SCHEMA IF NOT EXISTS raw")
con.execute("""
CREATE TABLE IF NOT EXISTS raw.day19_items (
id INTEGER, name VARCHAR, ts BIGINT, fetched_at TIMESTAMP
)
""")
total = 0
async with httpx.AsyncClient() as client:
session = PoliteSession(UA, qps=3.0)
async for page in paginate_offset(
client, "http://127.0.0.1:8766/api/offset", limit=50, session=session,
):
if not page.items:
break
now = datetime.now()
con.executemany(
"INSERT INTO raw.day19_items VALUES (?, ?, ?, ?)",
[(it["id"], it["name"], it["ts"], now) for it in page.items],
)
total += len(page.items)
print(f"累計 {total} 筆")
return total
if __name__ == "__main__":
n = asyncio.run(collect())
print(f"完成:raw.day19_items 共 {n} 筆")
這段用 paginate_offset 拿頁,再用 executemany 批次寫進 DuckDB。每次寫入都加上一個 fetched_at 時間戳,這是 Day 20 增量載入的關鍵欄位:明天我們會用這個欄位做「水位標記」,只抓「上次沒抓過」的資料。注意這裡用了「逐頁寫」而非「全部抓完再寫」:好處是中途失敗不會丟掉已經抓的;壞處是每頁一次 commit 在大量資料時較慢,Day 22 會用 transaction 包起來。
第五步:處理「429 Too Many Requests」這個 API 採集最常碰到的限流訊號:
"""pipelines/_rate_limit.py:處理 429 Retry-After 的禮貌重試。"""
import asyncio
import httpx
from datetime import datetime
async def fetch_with_429(client: httpx.AsyncClient, url: str, max_retries: int = 5) -> dict:
for attempt in range(max_retries):
r = await client.get(url)
if r.status_code != 429:
r.raise_for_status()
return r.json()
retry_after = r.headers.get("Retry-After")
if retry_after is None:
wait = 2 ** attempt
else:
try:
wait = float(retry_after)
except ValueError:
wait = 2 ** attempt
print(f"收到 429,等 {wait} 秒後重試(第 {attempt + 1}/{max_retries} 次)")
await asyncio.sleep(wait)
raise RuntimeError(f"超過 {max_retries} 次仍 429")
if __name__ == "__main__":
async def demo() -> None:
async with httpx.AsyncClient() as client:
data = await fetch_with_429(client, "http://127.0.0.1:8766/api/slow")
print(f"成功抓到 {len(data['items'])} 筆")
asyncio.run(demo())
這段是 Day 17 禮貌套件的延伸,專門處理 429。HTTP 429 通常伴隨 Retry-After header,這是「對方明確告訴你等多久再試」的訊號。Retry-After 可能是秒數(HTTP/1.1)或 HTTP 日期(HTTP/2),這裡用 try/except 包起來,無法解析時退回指數退避。2 ** attempt 是指數退避的核心:第 1 次等 1 秒、第 2 次等 2 秒、第 3 次等 4 秒;這對「短暫爆量」的狀況很有用,但對「持續爆量」就要靠「降低 QPS」根本解決。
第六步:做一個「分頁品質稽核」的 SQL 範本,抓完之後跑一次確認資料完整性:
"""pipelines/day19_audit.py:抓完之後的稽核 SQL。"""
import duckdb
con = duckdb.connect("warehouse/de-journey.duckdb")
print("=== 抓取摘要 ===")
print(con.execute("""
SELECT
COUNT(*) AS n,
COUNT(DISTINCT id) AS unique_ids,
MIN(id) AS min_id,
MAX(id) AS max_id,
MIN(fetched_at) AS first_fetch,
MAX(fetched_at) AS last_fetch
FROM raw.day19_items
""").fetchdf().to_string(index=False))
print("=== 重複 ID 檢查 ===")
dup = con.execute("""
SELECT id, COUNT(*) AS n FROM raw.day19_items
GROUP BY id HAVING COUNT(*) > 1
""").fetchall()
print(f"重複 ID 數:{len(dup)}(0 表示無重複)")
print("=== 連續性檢查 ===")
gaps = con.execute("""
WITH ids AS (SELECT DISTINCT id FROM raw.day19_items ORDER BY id),
diffs AS (SELECT id, id - LAG(id) OVER (ORDER BY id) AS gap FROM ids)
SELECT COUNT(*) AS gaps FROM diffs WHERE gap > 1
""").fetchone()[0]
print(f"ID 跳號數:{gaps}(0 表示連續)")
這段是 Day 18 稽核 SQL 的進階版。除了筆數與時間,還檢查「重複 ID 數」與「ID 跳號數」。重複 ID 是「抓了兩次以上」的訊號,常見原因有 offset 沒遞增、cursor 沒更新;ID 跳號是「資料被刪除或漏抓」的訊號,可能原因有 API 端資料被刪、分頁跨越刪除時間點。這兩個指標是分頁處理器的「健康檢查」,寫進排程後每天跑可以提早發現問題。實際結果依資料集而定,示範資料 1 到 250 應該無重複、無跳號。
第七步:寫一份「分頁處理器使用說明」當成交接文件的一部分:
"""pipelines/_pager_README.py:用程式碼輸出分頁處理器的使用說明。"""
import json
DOC = {
"模組": "_pager.py",
"提供": ["paginate_offset", "paginate_cursor"],
"設計": "AsyncIterator[Page] + 注入 PoliteSession",
"使用範本": """
async with httpx.AsyncClient() as client:
session = PoliteSession(UA, qps=3.0)
async for page in paginate_offset(client, URL, limit=100, session=session):
# 寫 DuckDB、寫檔、通知 Slack 等
process(page.items)
""",
"注意事項": [
"offset-based 對深頁面慢,cursor-based 對無限流友善",
"429 要讀 Retry-After header",
"重複 ID 表示分頁邏輯錯誤,跳號表示資料端變動",
"水位標記(watermark)由 Day 20 接管",
],
}
if __name__ == "__main__":
print(json.dumps(DOC, ensure_ascii=False, indent=2))
這份「可執行的文件」說明分頁處理器的設計與使用方式:它接受 httpx.AsyncClient 與 PoliteSession,吐出 AsyncIterator[Page],讓呼叫端自由決定處理方式。Day 40 的「文件化與交接」篇章會把這套概念延伸到完整的 README 與交接清單。執行這個檔案會印出 JSON 格式的說明,方便日後對照與覆審。
另一個實務上常見的延伸是「多分頁來源並行抓取」。一支完整的採集管線常會同時從好幾個 API 拿資料:交通部的即時車流、環保署的空氣品質、經濟部的水情。每個來源都有自己的節奏與限制,用 asyncio.gather() 搭配 PoliteSession 可以一次啟動多支抓取器。注意每個來源要帶自己的 User-Agent 與節流參數,避免「同一個 IP 對同一個主機用不同來源同時打」,那會被視為單一攻擊者。
實作上我們會在 Day 30 的端到端管線用一個 Collector class 把多個分頁來源收進來,並在每個來源結束後把摘要寫進 raw._collect_status。這套模式的精神是「每個來源獨立成功或失敗,不要一個失敗把全部拖下水」。Day 22 的失敗處理會展開這個觀念,今天先在心裡放個位置。
最後,分頁處理器的設計還有一個常被忽略的小細節:把 limit 與 offset(或 cursor)寫進「環境變數」而非硬編碼。這樣不同環境(開發、測試、正式)可以共用同一份程式碼,但用各自的參數。今天的範例把這些值寫在函式參數裡,算是最簡單的做法;Day 30 的端到端章節會引入 pydantic-settings 或 dynaconf 來管環境變數。
常見錯誤與踩雷
第一個雷是「忘了帶 limit 參數」。多數 API 的預設 limit 是 20 或 50,如果你只傳 offset 沒傳 limit,分頁會很慢、單頁資料太少。對應排查方向:先看 API 文件,確認預設 limit 與最大 limit;如果文件寫得不清楚,用「帶 limit=100 與不帶 limit」各測一次,比對筆數差異。
第二個雷是「offset 算錯」。例如 offset 從 0 開始,但下一頁應該從 limit 開始;或者你以為 offset 100 表示「從第 100 筆開始」,但 API 把它當成「跳過 100 筆」。對應排查方向:第一次跑之前先單頁測試,確認回傳的「下一頁 offset」是否等於「這一頁 offset + limit」;若是負值或亂碼,先停下來讀文件,不要硬跑。
第三個雷是「cursor 是 base64 編碼但你把它當純字串」。很多 API 的 cursor 是 base64 或 JSON 編碼過的「不透明字串」,你不能直接讀出裡面的數字。對應做法:把 cursor 當成「黑盒子」處理,不要嘗試解析或修改它;只要原封不動傳給下一頁就好。如果真的需要 debug,可以 base64 decode 看裡面的 JSON 結構(但通常沒有必要)。
第四個雷是「把禮貌節流關掉只為求快」。很多工程師在壓測時為了「把 10 萬筆資料塞進一次跑」,會把 QPS 拉到 20、50;這時對方伺服器可能直接 ban 你的 IP。對應排查方向:分頁處理器一定要保留禮貌節流,把「抓得多快」交給「加大 limit」而非「提高 QPS」。一次拿 500 筆比每秒打 5 次拿 100 筆更友善。Day 22 的失敗處理會再展開這個議題。
第五個雷是「分頁跨越時區或 DST」。時序資料若用「上次最大時間戳」當 cursor,碰到時區切換或日光節約時間可能會漏抓或重抓。對應做法:永遠用 UTC 儲存時間,cursor 內部也用 UTC;顯示給使用者看時再轉當地時區。DuckDB 的 TIMESTAMP 不帶時區,要存帶時區的時間用 TIMESTAMPTZ,這個觀念會在 Day 23 的維度建模篇章再展開。
效能與實務提醒
分頁處理器的效能瓶頸通常在「每頁的固定成本」:DNS 查詢、TCP 握手、TLS 驗證。如果每頁只取 50 筆,這些固定成本佔比會很高。對應做法:把 limit 拉到對方允許的最大(多數 API 是 100 或 500),用「減少頁數」換「每頁成本被稀釋」。另一個做法是用 HTTP/2(httpx 預設支援)與 keep-alive,讓同一個連線可以處理多頁。
另一個常見取捨是「同步 vs 非同步」。同步(requests)寫起來簡單,但單執行緒;非同步(httpx + asyncio)可並行,但 debug 較難。今天的範例全用非同步,是因為 Day 21 的排程、Day 22 的失敗處理都需要並行抓取多個來源的能力。建議從現在開始熟悉 asyncio,但不要一開始就寫太複雜——先用同步版驗證邏輯,再改成非同步版。
最後,「增量抓取」是 API 採集的核心能力之一。我們在第四步已經埋了 fetched_at 欄位,但完整的「水位標記」設計要等 Day 20 才會展開。今天的範例可以做為暖身:抓一次、抓兩次,比較 fetched_at 與 id 的差異,你就會看到「水位」的概念。今天做完之後,記得 pipelines/_pager.py 與 pipelines/day19_collect.py 都存起來,明天會拿來接水位標記。
小結
今天把 Day 17 與 Day 18 的能力擴充到「API 分頁」這層。我們了解了 offset-based 與 cursor-based 兩種主流分頁風格,寫了一個通用分頁處理器,處理了 429 限流,做了稽核 SQL 與交接文件。這套組合是 API 採集的標準配備,後續 Day 20(增量載入)、Day 30(端到端管線)都會反覆用到。把今天建立的關鍵詞整理進筆記本:offset、cursor、page iterator、429 Retry-After、指數退避、限流、配額。這些詞會在 Day 22 的失敗處理篇章再延伸為「重試佇列」與「斷點續跑」,在那之前先把分頁這層打穩。
結語
今天我們把「抓整包」推進到「分頁抓 API」,這是資料工程最常碰到的場景。但分頁只解決了「如何拿資料」,沒有解決「如何只拿新資料」。明天 Day 20 我們會把這個關鍵問題講清楚:水位標記(watermark)是什麼、為什麼它和冪等(idempotent)是設計增量載入的兩大支柱、如何在 DuckDB 裡實作水位表、如何驗證你的抓取真的冪等。這套概念會貫穿 Day 21 的每日排程與 Day 30 的端到端管線,是整個系列的工程核心之一。
明天,我們會在分頁處理器的基礎上,加上水位標記與冪等保證,把採集從「一次性」變成「每天可重複」。
延伸資源
- httpx 0.28 文件(2025):
https://www.python-httpx.org/,AsyncClient與 keep-alive 的標準用法。 - GitHub REST API 文件(2025):
https://docs.github.com/en/rest,真實世界中 link-based pagination 的最佳範例,Linkheader 給下一頁 URL。 - MDN HTTP 429(2025):
https://developer.mozilla.org/en-US/docs/Web/HTTP/Status/429,Too Many Requests 與Retry-Afterheader 的權威定義。 - DuckDB 1.4 INSERT 文件(2025):
https://duckdb.org/docs/sql/statements/insert.html,executemany批次插入的效能說明。 - 政府資料開放平臺 API(2025):
https://data.gov.tw/api/v2/rest,CKAN 風格的真實 API,依「政府資料開放授權條款第 1 版」可商用。
留言
張貼留言