DE Day 30 端到端管線(一):採集與落地
執行需求:CPU 可跑。今天是端到端管線系列的第一篇,我們要把前 29 篇學到的東西拼成一條可運行的管線,從「政府資料開放平台的真實資料集」一路落地到 DuckDB + Parquet,後續四天會繼續把轉換、建模、品質、監控、通知接上。本篇所有範例都在一般筆電的 CPU 上執行,不依賴 GPU。讀完這篇,你會有一個可以每天定時跑一次的管線雛形:把 data.gov.tw 的資料下載回來、保留原始檔、寫入 DuckDB、轉成 Parquet 分區落檔。執行前請先啟用 Day 2 建立的 de-journey/ 環境,並安裝本篇會用到的 httpx==0.28、tenacity==9.0、duckdb==1.4.1、pyarrow==18.0.0。
引言
前 29 篇我們把 SQL、DuckDB、Polars、pandas、爬蟲、開放資料、增量載入、排程、Airflow 等主題一個一個拆開來看。今天開始,我們要把這些零件組回一條真實的管線。系列會以「政府資料開放平台」的一個真實資料集為來源,沿著「採集 → 落地 → 轉換 → 建模 → 品質 → 監控 → 通知」的順序,每天往前推一步,最後在 Day 35 跑完整條管線、在 Day 36 與 Day 37 把排程器接上雲端。整套管線共用同一份 DuckDB 檔案、同一組資料表結構、同一組品質規則,這樣日後要改一個欄位或加一條規則時,改一處就全套生效。
本篇鎖定「採集與落地」這兩步:採集是「把資料從來源端抓回來」,落地是「把原始資料與中繼結果保存到本機」。我們選擇的資料集是經濟部在公司及商業設立登記一站式網站(位於 data.gov.tw)所提供的「公司登記資料」與「公司變更登記資料」兩個資料集,授權為「政府資料開放授權條款第 1 版」。這兩個資料集對應到「公司基本維度表」與「公司變更事實表」,正好是我們前面 Day 23–Day 26 講的維度建模的實戰標的。今天的目標很單純:寫出一個腳本,每天把這兩個資料集的最新版本拉回本機、寫到 DuckDB、再匯出成 Parquet 分區檔;明天的 Day 31 會把轉換與建模接到今天落地的資料上。
管線的目錄沿用 Day 1 的規劃:data/ 放原始檔(落地層)、warehouse/de-journey.duckdb 放 DuckDB 倉儲、pipelines/ 放管線腳本、logs/ 放執行紀錄。今天會新增兩個檔案:pipelines/ingest.py 負責採集與落地,pipelines/common.py 放共用設定(路徑、資料集 ID、欄位定義)。後續四天會在 pipelines/ 加更多腳本,但 common.py 與 warehouse/ 結構不會動,這是「同一份管線設定」的最小保證。
採集與落地的核心觀念
採集與落地是管線最前端的一步,也是最容易出包的一步。出包的模式大致有四種:來源端 5xx、網路瞬斷、檔案編碼錯誤、磁碟空間不足。實務上我們會用「重試、編碼宣告、磁碟檢查、雜湊驗證」四招來擋掉大部分狀況。重試的關鍵是「指數退避」:第一次失敗等 2 秒、第二次等 4 秒、第三次等 8 秒,最多三次。這樣能消化大多數的瞬時錯誤,又不會在真正的故障上耗太久。編碼的部分,政府資料平台常見的編碼是 UTF-8 與 Big5 混用;下載後先用 chardet 偵測,再決定讀檔方式。
落地的關鍵是「保留原始檔 + 寫入倉儲」兩個動作都要做。保留原始檔是為了「日後可以重跑」:當你發現轉換邏輯有 bug、或是政府平台改欄位定義時,可以從原始檔重新處理,而不用再回去抓一次。寫入倉儲則是給後續轉換與分析使用。我們用 DuckDB 的 read_csv_auto() 與 COPY ... TO 把 CSV 直接讀進來、再匯出成 Parquet。Parquet 的好處是「欄式壓縮 + 讀取時只挑需要的欄位」,對於幾十萬到幾百萬列的政府開放資料特別划算。
分區(partition)是另一個重要觀念。每天的資料會存到類似 data/company_basic/dt=2025-12-06/source=moea_basic.csv 的路徑,dt 是採集日、source 是資料來源簡稱。DuckDB 的 read_parquet('data/company_basic/dt=*/source=moea_basic/*.parquet', hive_partitioning=true) 可以用 glob 一次讀回所有分區,這是 Day 31 轉換時會用到的小技巧。分區的副作用是檔案數會爆炸,所以我們建議每季或每年合併一次(Day 35 的效能章節會講到合併策略)。
共通設定:管線常數與資料集約定
為了讓 Day 30 到 Day 35 共用同一份設定,我們把所有的路徑、資料集 ID、欄位定義、與品質門檻集中放在 pipelines/common.py。這個檔案今天會先建立最關鍵的三個常數:PROJECT_ROOT、DATASETS 與 DUCKDB_PATH。後續幾天會在同一個檔案加上品質規則清單、通知設定與排程設定。請把這份設定當成「管線的單一事實來源(single source of truth)」:改一處,全套生效。
"""de-journey/pipelines/common.py:Day 30-35 共用的管線常數。
請把這份檔案視為「管線的單一事實來源」。任何路徑、資料集 ID、欄位名稱
或品質門檻的變動都先改這裡,再讓其他腳本讀取。
"""
from pathlib import Path
# 工作目錄結構沿用 Day 1 的規劃
PROJECT_ROOT = Path(__file__).resolve().parents[1]
DATA_DIR = PROJECT_ROOT / "data"
WAREHOUSE_DIR = PROJECT_ROOT / "warehouse"
LOGS_DIR = PROJECT_ROOT / "logs"
DUCKDB_PATH = WAREHOUSE_DIR / "de-journey.duckdb"
# 兩個資料集都來自 data.gov.tw,授權為政府資料開放授權條款第 1 版
# 檔案 URL 以 data.gov.tw 平台公開的檢視頁為入口;實際檔案連結需到平台查詢
DATASETS = {
"company_basic": {
"title": "公司登記基本資料",
"source": "moea_basic",
# 對應的 DuckDB 維度表
"table": "raw.company_basic",
# 重要的欄位(其他欄位仍會保留,但品質檢查只看這些)
"key_columns": ["uniform_no", "company_name"],
"partition_prefix": "company_basic",
},
"company_change": {
"title": "公司變更登記資料",
"source": "moea_change",
"table": "raw.company_change",
"key_columns": ["uniform_no", "change_date", "change_item"],
"partition_prefix": "company_change",
},
}
# 下載與重試的預設值(給 httpx + tenacity 用)
HTTP_TIMEOUT_SEC = 30
RETRY_ATTEMPTS = 3
RETRY_BACKOFF_SEC = 2.0
這份設定檔不大,但它定義了三件後續每天都會用到的事:第一,DUCKDB_PATH 是倉儲檔案的絕對路徑,所有 SQL 連線都指向它,這樣 Day 30 寫進去的資料、Day 32 跑品質檢查、Day 34 在 Streamlit 讀取時都是同一份檔案;第二,DATASETS 用字典統一管理「資料集名稱 → 對應的 DuckDB 表、欄位清單、落地路徑前綴」,將來加第三個、第四個資料集只要在這裡加條目;第三,HTTP_TIMEOUT_SEC 與重試設定統一在這裡管理,避免每個腳本各寫一份「我自己的 timeout」。
先在 pipelines/__init__.py 留一個空檔,讓上面的 import 與後續的 python -m pipelines.ingest 都能用 package 形式呼叫:
# de-journey/pipelines/__init__.py:把 pipelines/ 標成 Python package
# 內容留空即可;所有子腳本會被以 python -m pipelines.<name> 形式執行。
接著建立倉儲與落地目錄。實務上你會希望 warehouse/ 與 data/ 在第一次跑管線前就存在:
mkdir -p warehouse logs
mkdir -p data/company_basic data/company_change
ls -1
# 輸出:
# data
# logs
# pipelines
# pyproject.toml
# warehouse
環境檢查腳本:每天開工前先跑一次,確保 Python 與核心套件版本都在預期範圍:
"""de-journey/check_env.py:開工前先跑這個,確認版本到位。"""
import sys
MIN_PYTHON = (3, 13)
EXPECTED = {
"duckdb": "1.4",
"polars": "1.3",
"pandas": "2.3",
"httpx": "0.28",
"tenacity": "9.0",
}
# Python 主版本檢查
if sys.version_info[:2] < MIN_PYTHON:
raise SystemExit(
f"Python 版本過舊:{sys.version.split()[0]},需要 {MIN_PYTHON[0]}.{MIN_PYTHON[1]}+"
)
for name, prefix in EXPECTED.items():
mod = __import__(name)
version = getattr(mod, "__version__", "unknown")
if not version.startswith(prefix):
raise SystemExit(f"{name} 版本 {version} 不符合預期 {prefix}.x")
print(f"{name}:{version}") # 輸出:httpx:0.28.1
這支腳本會在版本不符合預期時直接 raise,避免管線跑到一半才發現套件太舊。把這個腳本當成「每天開工的健康檢查」,搭配 Day 32 的品質檢查一起用。
完整實作:採集 + 落地到 DuckDB 與 Parquet
接下來把 pipelines/ingest.py 寫完整。這支腳本的工作流程是:對 DATASETS 中的每個資料集,先檢查磁碟空間、再用 httpx 下載到當天的落地目錄、接著用 DuckDB 讀進 raw.[name]、最後匯出成 Parquet 分區檔。整段約 100 行,可以直接貼到 pipelines/ingest.py。執行前需要:uv pip install httpx==0.28.1 tenacity==9.0.0 duckdb==1.4.1 pyarrow==18.0.0。
"""de-journey/pipelines/ingest.py:Day 30 採集與落地腳本。
對 DATASETS 裡的每個資料集執行:
1. 檢查磁碟空間
2. 下載原始檔到 data/<name>/dt=<today>/
3. 用 DuckDB 讀進 raw.<name>
4. 匯出 Parquet 到 data/<name>/dt=<today>/parquet/
執行:
python -m pipelines.ingest
"""
import logging
import os
import shutil
import sys
from datetime import date
from pathlib import Path
import duckdb
import httpx
from tenacity import (
retry,
stop_after_attempt,
wait_exponential,
retry_if_exception_type,
)
from pipelines.common import (
DATA_DIR,
DUCKDB_PATH,
HTTP_TIMEOUT_SEC,
LOGS_DIR,
DATASETS,
RETRY_ATTEMPTS,
RETRY_BACKOFF_SEC,
)
# 把 logging 導到 logs/ingest-YYYY-MM-DD.log
LOGS_DIR.mkdir(parents=True, exist_ok=True)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
handlers=[
logging.FileHandler(LOGS_DIR / f"ingest-{date.today():%Y-%m-%d}.log"),
logging.StreamHandler(sys.stdout),
],
)
log = logging.getLogger("ingest")
# 預留 50 MB 緩衝,避免磁碟滿載
MIN_FREE_BYTES = 50 * 1024 * 1024
def ensure_disk_space(path: Path, need_bytes: int = MIN_FREE_BYTES) -> None:
"""檢查落地目錄所在磁碟的可用空間。"""
usage = shutil.disk_usage(path)
free = usage.free
if free < need_bytes:
raise RuntimeError(
f"磁碟可用空間不足:剩餘 {free} bytes,門檻 {need_bytes} bytes"
)
@retry(
stop=stop_after_attempt(RETRY_ATTEMPTS),
wait=wait_exponential(multiplier=RETRY_BACKOFF_SEC, min=2, max=30),
retry=retry_if_exception_type((httpx.HTTPError, ConnectionError)),
reraise=True,
)
def download_to(url: str, dest: Path) -> int:
"""下載 URL 到 dest,回傳檔案大小(bytes)。
對瞬時錯誤(5xx、連線重置)做指數退避重試。
"""
dest.parent.mkdir(parents=True, exist_ok=True)
log.info("下載中:%s", dest)
with httpx.Client(timeout=HTTP_TIMEOUT_SEC, follow_redirects=True) as client:
with client.stream("GET", url) as response:
response.raise_for_status()
with dest.open("wb") as fh:
for chunk in response.iter_bytes():
fh.write(chunk)
size = dest.stat().st_size
log.info("下載完成:%s(%d bytes)", dest, size)
return size
def land_to_duckdb(csv_path: Path, table: str) -> int:
"""把 CSV 讀進 DuckDB 的指定表,回傳資料筆數。
用 CREATE OR REPLACE TABLE 把當天的快照整批寫入。
"""
con = duckdb.connect(str(DUCKDB_PATH))
try:
con.execute("CREATE SCHEMA IF NOT EXISTS raw")
# 政府開放資料常見編碼為 UTF-8;少數為 Big5,read_csv_auto 會自動判斷
con.execute(
f"CREATE OR REPLACE TABLE {table} AS "
f"SELECT * FROM read_csv_auto(?, header=true, sample_size=-1)",
[str(csv_path)],
)
n = con.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
log.info("已寫入 %s:%d 筆", table, n)
return n
finally:
con.close()
def export_to_parquet(csv_path: Path, parquet_dir: Path) -> int:
"""把 CSV 轉成 Parquet 分區檔,回傳檔案大小。"""
parquet_dir.mkdir(parents=True, exist_ok=True)
con = duckdb.connect(str(DUCKDB_PATH))
try:
con.execute(
"COPY (SELECT * FROM read_csv_auto(?, header=true)) TO ? (FORMAT PARQUET, COMPRESSION zstd)",
[str(csv_path), str(parquet_dir / "part-000.parquet")],
)
finally:
con.close()
return sum(p.stat().st_size for p in parquet_dir.glob("*.parquet"))
def run_one(name: str, cfg: dict, today: date) -> None:
"""執行單一資料集的採集與落地。"""
url = os.environ.get(f"DATASET_URL_{name.upper()}")
if not url:
raise RuntimeError(
f"未提供 {name} 的下載 URL,請設定環境變數 DATASET_URL_{name.upper()}"
)
land_dir = DATA_DIR / cfg["partition_prefix"] / f"dt={today:%Y-%m-%d}"
csv_path = land_dir / f"{cfg['source']}.csv"
parquet_dir = land_dir / "parquet"
ensure_disk_space(land_dir)
download_to(url, csv_path)
n_rows = land_to_duckdb(csv_path, cfg["table"])
n_bytes = export_to_parquet(csv_path, parquet_dir)
log.info("[%s] 完成:%d 筆、parquet=%d bytes", name, n_rows, n_bytes)
def main() -> int:
today = date.today()
for name, cfg in DATASETS.items():
log.info("=== %s 開始 ===", name)
run_one(name, cfg, today)
log.info("=== 全部完成 ===")
return 0
if __name__ == "__main__":
sys.exit(main())
這段腳本的結構對應到前面講的「重試、編碼、磁碟、雜湊」四招。ensure_disk_space() 用 shutil.disk_usage() 檢查落地目錄所在磁碟的可用空間,低於 50 MB 就直接 raise;這比「下載到一半才發現磁碟滿」要好得多。download_to() 用 httpx.stream 做流式下載(避免大檔一次讀進記憶體),再用 tenacity 的 wait_exponential 設定指數退避(2 秒、4 秒、8 秒、最多 30 秒),並對 httpx.HTTPError 與 ConnectionError 重試。land_to_duckdb() 用 CREATE OR REPLACE TABLE 把當天的快照寫進 raw.[name],這是「落地層」的設計:當天看到的資料長什麼樣、就完整保存下來,後續轉換層從這份資料讀,不會被中途改欄位影響。
export_to_parquet() 是另一個重要動作:把同一份資料匯出成 Parquet 分區檔,用 zstd 壓縮。為什麼要同時存 CSV 與 Parquet?CSV 是人類可讀、方便除錯;Parquet 是分析友善、壓縮比高、欄位讀取快。Day 35 的效能章節會用 EXPLAIN ANALYZE 比較兩者的查詢時間差異。在實務上我們會建議「CSV 留 30 天當觀察期、Parquet 留長期」。
腳本最後的 main() 對 DATASETS 字典裡的每個資料集依序執行。URL 用環境變數注入(DATASET_URL_COMPANY_BASIC、DATASET_URL_COMPANY_CHANGE),這是「12-factor app」的標準做法:設定與程式碼分離,方便部署到不同環境時只改環境變數、不改程式碼。實際的 URL 由排程器(Day 36 / Day 37)或本機的 .env 檔提供;data.gov.tw 平台的檢視頁會列出每個資料集當下的下載連結,每次更新可能會換,建議把 URL 維護在一個獨立的小檔(例如 pipelines/secrets.example.env)並由 git 忽略。
驗證落地結果
跑完採集與落地後,我們要立刻驗證資料有沒有寫對。Day 14 介紹過的 DuckDB 內建 SQL 很適合做這件事:
"""de-journey/pipelines/verify_landing.py:跑完 ingest.py 後驗證落地結果。"""
import duckdb
from pipelines.common import DATASETS, DUCKDB_PATH
con = duckdb.connect(str(DUCKDB_PATH))
for name, cfg in DATASETS.items():
n_rows = con.execute(f"SELECT COUNT(*) FROM {cfg['table']}").fetchone()[0]
n_cols = con.execute(
f"SELECT COUNT(*) FROM information_schema.columns "
f"WHERE table_schema='raw' AND table_name='{cfg['table'].split('.')[1]}'"
).fetchone()[0]
head = con.execute(f"SELECT * FROM {cfg['table']} LIMIT 3").fetchall()
print(f"{name}: {n_rows} 筆、{n_cols} 欄")
for row in head:
print(f" {row}")
con.close()
# 輸出(實際數字會依當日下載而略有不同):
# company_basic: 712345 筆、18 欄
# ('10458575', '○○有限公司', ...)
# company_change: 89102 筆、12 欄
# ('10458575', '2024-08-15', '資本額變更', ...)
這支驗證腳本做了三件事:第一,確認每個資料表的總筆數(公司基本資料通常落在 60 萬至 80 萬筆、變更登記落在 8 萬至 12 萬筆);第二,確認欄位數量與政府平台公布的規格一致;第三,印出前三列樣本,肉眼檢查欄位有沒有錯位(例如公司名稱跑到地址欄)。如果筆數或欄位數與預期差太多,就是下載或讀檔環節出問題,需要回頭查日誌。
如果要把落地檔案的指紋記下來方便日後比對,可以用 DuckDB 的 hash() 函式對每個分區算 SHA-256 並寫回一份 metadata 表:
"""de-journey/pipelines/hash_landing.py:對當天落地的 CSV/Parquet 算指紋。"""
import duckdb
from pathlib import Path
from pipelines.common import DATA_DIR, DATASETS, DUCKDB_PATH
con = duckdb.connect(str(DUCKDB_PATH))
con.execute("CREATE SCHEMA IF NOT EXISTS meta")
con.execute(
"CREATE TABLE IF NOT EXISTS meta.file_fingerprint ("
" dataset VARCHAR, dt DATE, path VARCHAR, sha256 VARCHAR, "
" bytes BIGINT, PRIMARY KEY (dataset, dt, path))"
)
today = Path.cwd().name # 簡化版;實務上從環境變數或 sys.argv 讀
for name, cfg in DATASETS.items():
src = DATA_DIR / cfg["partition_prefix"]
if not src.exists():
continue
for csv in src.glob("**/*.csv"):
sha = con.execute(
"SELECT hash(?, 'sha256') AS h", [str(csv)]
).fetchone()[0]
size = csv.stat().st_size
con.execute(
"INSERT OR REPLACE INTO meta.file_fingerprint "
"VALUES (?, current_date, ?, ?, ?)",
[name, str(csv), sha, size],
)
print(f"{csv.name} -> {sha[:16]}... ({size} bytes)") # 輸出:sha256 前 16 字
con.close()
這份指紋表會在後續 Day 32 品質檢查時派上用場:當我們發現今天的資料筆數異常變少時,可以對照昨天的指紋確認「是不是檔案被截斷」、「還是政府平台真的更新了內容」。
常見錯誤與踩雷
錯誤一:忘了設定資料集 URL,腳本第一行就 raise。常見症狀:RuntimeError: 未提供 company_basic 的下載 URL。對應排查方向:請確認環境變數 DATASET_URL_COMPANY_BASIC 與 DATASET_URL_COMPANY_CHANGE 都有設定;或是在 .env 檔中讀取後用 export 載入到環境。另一個解法是把 URL 直接寫到 DATASETS 字典(拿掉環境變數),但這會讓設定與程式碼混在一起,不建議在正式環境用。
錯誤二:Big5 編碼的 CSV 被誤判成 UTF-8。常見症狀:讀進 DuckDB 後中文欄位出現亂碼。對應排查方向:用 chardet 或 charset-normalizer 偵測編碼後,在 read_csv_auto 加上 encoding='big5' 參數。對應到本篇的 land_to_duckdb 函式,可以改成:read_csv_auto(?, header=true, encoding=?, sample_size=-1) 並在呼叫端偵測編碼。
錯誤三:磁碟空間檢查通過但下載到一半滿了。常見症狀:httpx 拋出 OSError: No space left on device。對應排查方向:ensure_disk_space 的 50 MB 預留只是「下載前的快照」,如果中間有其他 process 寫入大量檔案,空間還是會被吃光。解法是在 download_to() 內每寫 1 MB 就重新檢查一次剩餘空間,或在排程時把 de-journey/ 放在專屬磁碟(避免被瀏覽器快取、log 檔塞滿)。
錯誤四:CREATE OR REPLACE TABLE 把昨天的快照蓋掉了。這個不是 bug、是設計:我們把 raw.[name] 設計成「每天的最新快照」,歷史資料靠 Parquet 分區檔保存。如果你要保留每天的歷史快照到 DuckDB 裡,可以改成 CREATE TABLE IF NOT EXISTS 加上日期欄位,或是用 INSERT INTO 累加;但要注意「重跑同一日」時不要重複寫入(這是 Day 31 會講的冪等議題)。
錯誤五:httpx 預設沒有重試,要記得用 tenacity 顯式加上。常見症狀:第一次遇到 503 就整支腳本 fail。對應排查方向:請確認 download_to() 上面有 @retry(...) 裝飾子,並且 retry_if_exception_type 包含 httpx.HTTPError 與 ConnectionError。另一個常見問題是忘了設定 reraise=True,導致重試失敗後被 silently swallow;這個 flag 一定要打開,否則你會以為下載成功了,實際上檔案是空的。
效能與實務提醒
採集與落地的效能瓶頸通常在「網路」與「磁碟 I/O」,不在 Python 本身。以經濟部公司登記資料為例,CSV 檔案大小約 80–120 MB,在家用網路下載時間約 10–20 秒;DuckDB 寫入 70 萬筆約 3–5 秒;Parquet 匯出約 1–2 秒。整支腳本從開始到結束應該在 1 分鐘以內完成。若你發現下載時間動輒超過 5 分鐘,先檢查網路;如果網路沒問題,那就是平台當下流量大或檔案變大。
實務上有兩個取捨要記得。第一個是「CSV 留多久」。政府開放資料的 CSV 通常 30 天就會被新版本蓋掉,所以留太久沒意義;建議 30 天滾動保留,用 find data/ -mtime +30 -delete 排程清掉。第二個是「Parquet 壓縮等級」。zstd 是速度與壓縮比的最佳平衡點(比 gzip 快、比 snappy 小);如果你要更小檔案可以用 brotli,但 CPU 時間會增加 2–3 倍。一般情況下 zstd level 3 就夠了,不要為了 5% 的空間節省犧牲 30% 的 CPU 時間。
另一個工程建議:把 ingest.py 與 verify_landing.py 綁在一起跑(python -m pipelines.ingest && python -m pipelines.verify_landing)。驗證腳本的執行時間只有幾秒,但它能在第一時間告訴你「落地成功且資料合理」。如果哪天你看到驗證腳本的欄位數量從 18 變成 17,就代表政府平台改了欄位定義,這時候要進到轉換層(Day 31)做對應調整。
小結
今天把端到端管線的第一塊拼上去:採集與落地。我們建立了 pipelines/common.py 作為管線的單一事實來源,定義了兩個資料集(公司登記基本、公司變更登記)的 DuckDB 表與 Parquet 落地路徑;用 ingest.py 把 HTTP 下載、DuckDB 寫入、Parquet 匯出三件事串起來,加入磁碟空間檢查、指數退避重試、logging 到檔案等工程細節;並用 verify_landing.py 做落地後的快速驗證。重點回顧:第一,採集端要處理「瞬時錯誤」與「編碼錯誤」,這兩個是政府開放資料最常見的失敗模式;第二,落地層要同時保留 CSV 與 Parquet,CSV 給人看、Parquet 給機器讀;第三,CREATE OR REPLACE TABLE 是「每天最新快照」的設計,歷史要靠 Parquet 分區檔找回來;第四,設定與程式碼分離(環境變數、.env 檔)是 12-factor 的基本要求;第五,落地後一定要驗證筆數與欄位數,否則 Day 31 的轉換會在錯的資料上跑。
明天 Day 31 會接著做「轉換與建模」:把今天落地的 raw.company_basic 與 raw.company_change 透過 SQL 轉成清洗後的 staging.company_basic_clean 與 staging.company_change_clean,再組成維度表 mart.dim_company 與事實表 mart.fact_company_change。轉換層會用到 Day 3–Day 7 的 SQL 技巧(視窗函式、CTE)、Day 13–Day 14 的清洗與品質觀念,以及 Day 23–Day 24 的維度建模框架。
結語
今天的重點是「把前 29 篇的零件組回一條可運行的管線」。我們沒有急著把所有功能一次寫完,而是先把「採集與落地」這最關鍵的一步做穩:URL 從環境變數讀、DuckDB 用快照式寫入、Parquet 用 zstd 壓縮分區、磁碟空間與網路錯誤都有基本保護。這套骨架在 Day 31–Day 35 不會變,只會加更多處理步驟與品質規則。
明天,我們會把今天落地的 raw 表進一步整理成 staging 與 mart 兩層:清洗、型別轉換、維度表與事實表的產製。我們會用一個可以整段執行的 SQL 範例,把「同一家公司的所有變更事件」壓平成一張事實表、把「公司基本資料」整理成一張維度表,並用 Day 23 介紹的 star schema 概念把它們接起來。這是 Day 32 品質檢查與 Day 34 Streamlit 儀表板的基礎。
延伸資源
- 政府資料開放平臺:
https://data.gov.tw/。本系列使用的「公司登記資料」與「公司變更登記資料」位於此平台,授權為「政府資料開放授權條款第 1 版」,使用時須保留出處標示。 - httpx 官方文件(2025,0.28 版):
https://www.python-httpx.org/。本篇的Client.stream()與follow_redirects用法以這份文件為準;如需更複雜的併發下載,可改用httpx.AsyncClient與asyncio.gather()。 - tenacity 官方文件(2025,9.0 版):
https://tenacity.readthedocs.io/。指數退避wait_exponential與條件重試retry_if_exception_type是本篇的核心設定。 - DuckDB CSV 讀取官方文件(2025,1.4 版):
https://duckdb.org/docs/stable/data/csv/reading。read_csv_auto()的編碼自動判斷、sample_size=-1全檔掃描選項,是本篇land_to_duckdb的依據。 - Parquet 官方格式說明(2025):
https://parquet.apache.org/docs/。本篇使用 zstd 壓縮、snappy 是另一個常見選項,差異在壓縮比與解壓速度的取捨。 - 12-Factor App:
https://12factor.net/config。本篇以環境變數注入資料集 URL,正是「設定與程式碼分離」的標準做法;Day 37 的 GitHub Actions 會用同一個原則管理 secrets。
留言
張貼留言