跳到主要內容

DE Day 30 端到端管線(一):採集與落地



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。

留言

這個網誌中的熱門文章

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 建構深度學習模型。 開發者與研究人員 :想更深入了...