跳到主要內容

DE Day 42 管線實作與排程

DE Day 42 管線實作與排程

執行需求:CPU 可跑。昨天把專案範圍、合約、模型定義清楚,今天要讓管線真的跑起來。我們會寫一支完整的每日管線(ingest → transform → test → notify),用 APScheduler 3.11 設定本機每日 03:00 (UTC+8) 的自動執行,並用 GitHub Actions 把同樣的流程部署到雲端做為備援。整個流程沿用 Day 41 建立的 project_config.py 與 aqi_hourly.yml,所有 Day 41-45 的程式碼引用同一份共用設定。讀完之後你應該能回答:每日管線的「成功完成」怎麼定義?APScheduler 與 GitHub Actions 怎麼分工?管線失敗時怎麼自動重試?

引言

管線設計有兩個常見的極端。第一個極端是「先把所有東西寫完再談排程」,結果跑了兩個月才發現某個步驟在排程環境失敗。第二個極端是「馬上把每個腳本都接到 cron」,結果缺乏監控、失敗沒人知道、排程亂成一團。我們採取第三條路線:先把 ingest、transform、test、notify 四個階段的單元腳本寫好並在命令列驗證,再接進 APScheduler 與 GitHub Actions。這樣每個階段的失敗訊息都是明確的,不會被排程器的「背景執行」模糊掉。

今天的設計原則是「單一入口、可重現、可監控」。單一入口指 pipelines/daily_run.py 這一支腳本負責一整天的所有工作,無論是手動跑、排程器跑、CI 跑都呼叫同一支腳本;可重現指同一個 --date 參數跑兩次會得到一樣的最終結果(冪等);可監控指每個階段都會寫 log、出問題會寫 breach 檔。

這篇文章會做四件事。第一件,寫 ingest 腳本:從環境部 CSV(或合成資料)抓取前一天的逐時資料,落地成 Parquet 與 DuckDB 表。第二件,補上 Day 41 還沒寫的 marts 模型(fct_aqi_hourly、dim_stations、dim_counties、mart_daily_summary)。第三件,把 Day 38 的合約檢查接進管線,自動跑 SLI 並寫 breach。第四件,用 APScheduler 與 GitHub Actions 兩種方式設定每日 03:00 排程。

共用設定:沿用 Day 41 的 project_config

今天的 ingest、transform、scheduler 都會 from project_config import MODELS, INDICATORS, SCHEDULE, ALERT。這代表如果未來要改模型名稱、改指標定義、改排程時間,只要改 project_config.py 一個檔案。Day 41 已經把 MODELS、INDICATORS、SCHEDULE 寫好,今天直接 import 即可:

"""de-journey/pipelines/daily_run.py:每日管線的單一入口。"""
from __future__ import annotations

import argparse
import logging
import sys
from datetime import datetime, timedelta, timezone
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parent))
from project_config import PROJECT, MODELS, INDICATORS, SCHEDULE, ALERT  # noqa: E402

# 1. 設定 logging:同時輸出到終端機與 logs/daily_run_<date>.log
log_dir = Path(__file__).resolve().parents[1] / "logs"
log_dir.mkdir(parents=True, exist_ok=True)
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
    handlers=[
        logging.StreamHandler(),
        logging.FileHandler(log_dir / "daily_run.log", encoding="utf-8"),
    ],
)
log = logging.getLogger("daily_run")

log.info("管線啟動,專案:%s", PROJECT["name"])
log.info("資料來源:%s", PROJECT["source"]["name"])
# 輸出(節錄):管線啟動,專案:aqi_pipeline
# 輸出(節錄):資料來源:行政院環境部(airtw.moenv.gov.tw)

這段把 logging 與共用設定都準備好。handlers 同時掛終端機與檔案,這樣手動跑時看得到輸出、排程器跑時有檔案可查。logging 的格式 %(asctime)s %(levelname)s %(message)s 是 Day 40 推薦的標準格式,方便後續用 grep 與 jq 處理。

第一步:ingest 腳本

ingest 階段的目標是「把前一天的逐時量測從環境部抓下來,寫進 DuckDB」。為了讓範例可重現,這個階段同時支援「真實 API/CSV 下載」與「合成資料」兩種模式,用 --use-synthetic 旗標切換。實務上你會把 --use-synthetic 拿掉,只留真實抓取。

"""de-journey/pipelines/ingest_aqi.py:從環境部抓取逐時 AQI 並寫進 DuckDB。"""
from __future__ import annotations

import argparse
import logging
import random
import sys
from datetime import datetime, timedelta, timezone
from pathlib import Path

import duckdb

sys.path.insert(0, str(Path(__file__).resolve().parent))
from project_config import PROJECT  # noqa: E402

log = logging.getLogger("ingest_aqi")
WAREHOUSE = Path(__file__).resolve().parents[1] / "warehouse" / "de-journey.duckdb"


def ingest_synthetic(target_date: datetime.date) -> int:
    """產生當天的合成 AQI 資料,覆寫到 DuckDB 的 aqi_raw schema。"""
    random.seed(target_date.toordinal())
    stations = ["466920", "466921", "466922", "466930", "466940",
                "467050", "467060", "467110", "467350", "467480"]
    rows = []
    for sid in stations:
        for hour in range(24):
            ts = datetime.combine(target_date, datetime.min.time(),
                                  tzinfo=timezone(timedelta(hours=8)))
            ts = ts.replace(hour=hour)
            base = random.gauss(50, 20)
            rows.append({
                "station_id": sid,
                "observed_at": ts,
                "aqi": max(0, min(500, int(base + random.gauss(0, 10)))),
                "pm25": max(0.0, min(500.0, round(base * 0.6 + random.gauss(0, 3), 1))),
                "county": random.choice(["臺北市", "新北市", "臺中市", "高雄市"]),
                "ingestion_at": datetime.now(timezone(timedelta(hours=8))),
            })

    con = duckdb.connect(str(WAREHOUSE))
    con.execute("CREATE SCHEMA IF NOT EXISTS aqi_raw")
    con.execute(
        "CREATE OR REPLACE TABLE aqi_raw.aqi_hourly AS "
        "SELECT * FROM (VALUES {})".format(
            ", ".join(["(CAST(? AS VARCHAR), CAST(? AS TIMESTAMP), "
                       "CAST(? AS INTEGER), CAST(? AS DOUBLE), "
                       "CAST(? AS VARCHAR), CAST(? AS TIMESTAMP))"] * len(rows))
        )
    )
    # 用 Polars 寫入更乾淨,但這裡為了單一檔案可跑,示範 insert 寫法
    import polars as pl
    con.execute("CREATE OR REPLACE TABLE aqi_raw.aqi_hourly AS SELECT * FROM "
                "((SELECT * FROM pl.DataFrame($df)) ".replace("$df", "df"))
    con.close()
    log.info("合成 ingest 完成:%d 列", len(rows))
    return len(rows)


def parse_args() -> argparse.Namespace:
    p = argparse.ArgumentParser()
    p.add_argument("--date", default=datetime.now(timezone(timedelta(hours=8)))
                   .date().isoformat())
    p.add_argument("--use-synthetic", action="store_true",
                   help="使用合成資料(示範用,真實部署請移除)")
    return p.parse_args()


if __name__ == "__main__":
    args = parse_args()
    target = datetime.fromisoformat(args.date).date()
    if args.use_synthetic:
        n = ingest_synthetic(target)
        print(f"已 ingest {n} 列(合成)")
    else:
        print("請從環境部 API 拉取真實資料;本範例僅示範合成模式")

這段把 ingest 腳本寫完。設計上有三個關鍵:第一,--use-synthetic 旗標讓範例可重現又保持真實部署的擴充性;第二,seed = target_date.toordinal() 讓同一天的合成資料固定下來,這樣跑兩次 --date 2025-12-09 會得到一模一樣的資料;第三,每列都帶 ingestion_at 時間戳,這對 Day 43 的 freshness 計算很重要。

我故意把 con.execute(...) 的寫法簡化,實際上更乾淨的寫法是用 Polars 的 pl.DataFrame 直接寫進 DuckDB:

# 更乾淨的 ingest 寫法:用 Polars 寫進 DuckDB(替換上面那段 execute)
import polars as pl
con = duckdb.connect(str(WAREHOUSE))
df = pl.DataFrame(rows)
con.execute("CREATE SCHEMA IF NOT EXISTS aqi_raw")
con.register("df_view", df)
con.execute("CREATE OR REPLACE TABLE aqi_raw.aqi_hourly AS SELECT * FROM df_view")
log.info("已寫入 %d 列到 aqi_raw.aqi_hourly", df.height)
con.close()
return df.height

這個版本更乾淨:先建 pl.DataFrame,再用 DuckDB 的 register 把 DataFrame 註冊成 DuckDB 視圖,最後 CREATE OR REPLACE TABLE 把視圖落地。DuckDB 對 register 的 Polars DataFrame 是零拷貝的,速度比 pandas 互轉快很多(Day 11 與 Day 12 比較過)。實務上你會選這個寫法。

輸出範例(--date 2025-12-09 --use-synthetic):

已 ingest 240 列(合成)

240 列 = 10 個測站 × 24 個小時。這個數字與真實環境部每天公布的「測站 × 小時」數量級一致(真實全台約 60+ 測站,一天約 1,500 列),範例用 10 個測站是為了讓 ingest 跑得快。實務上把 stations 換成真實清單即可。

第二步:補上 marts 模型

Day 41 寫了 staging 與 intermediate,今天補上四個 marts 模型。fct_aqi_hourly 是事實表,dim_stations 與 dim_counties 是維度表,mart_daily_summary 是聚合寬表。沿用 Day 23 的星狀模型設計,每張表的顆粒度都在檔案開頭明確註解。

-- de-journey/models/aqi_project/models/marts/fct_aqi_hourly.sql
{{ config(materialized='table') }}

-- 顆粒度 = 測站 × 小時
SELECT
    f.observed_at,
    f.station_id,
    s.station_name,
    s.county,
    s.operator,
    s.latitude,
    s.longitude,
    f.aqi,
    f.pm25,
    f.ingestion_at
FROM {{ ref('int_aqi_with_station') }} f
    LEFT JOIN {{ ref('dim_stations') }} s
        ON f.station_id = s.station_id

事實表的設計原則沿用 Day 23:顆粒度明確、量測欄位可加總(aqi、pm25 是數字、可加總)、維度表資訊透過 LEFT JOIN 補上。operator(營運單位,例如「行政院環境部」、「地方政府」)是這次新加的維度欄位,對 Day 44 儀表板的「依營運單位比較」有用。

-- de-journey/models/aqi_project/models/marts/dim_stations.sql
{{ config(materialized='table') }}

SELECT DISTINCT
    station_id,
    station_name,
    county,
    address,
    operator,
    latitude,
    longitude
FROM {{ ref('stg_stations') }}
WHERE station_id IS NOT NULL

DISTINCT 是這張維度表的關鍵:當 stg_stations 同一個測站有多列(例如運營單位變更、地址變更)時,DISTINCT 保證 dim_stations 每一個測站只有一列。實務上當你需要保留「測站變更歷史」時,要改用 Type 2 SCD(Day 23 提過),這裡先做最簡單的 Type 1。

-- de-journey/models/aqi_project/models/marts/dim_counties.sql
{{ config(materialized='table') }}

-- 全台 22 縣市的固定維度表(人工維護,每年更新一次)
SELECT * FROM (VALUES
    ('臺北市', '北部', 2646000, 272),
    ('新北市', '北部', 4034000, 2053),
    ('桃園市', '北部', 2306000, 1220),
    ('臺中市', '中部', 2842000, 2215),
    ('臺南市', '南部', 1866000, 2192),
    ('高雄市', '南部', 2752000, 2952),
    ('基隆市', '北部', 363000, 133),
    ('新竹市', '北部', 451000, 104),
    ('新竹縣', '北部', 580000, 1428),
    ('苗栗縣', '北部', 538000, 1820),
    ('彰化縣', '中部', 1245000, 1074),
    ('南投縣', '中部', 480000, 4106),
    ('雲林縣', '中部', 666000, 1290),
    ('嘉義市', '南部', 268000, 60),
    ('嘉義縣', '南部', 491000, 1903),
    ('屏東縣', '南部', 802000, 2776),
    ('宜蘭縣', '東部', 449000, 2144),
    ('花蓮縣', '東部', 318000, 4628),
    ('臺東縣', '東部', 213000, 3515),
    ('澎湖縣', '離島', 105000, 127),
    ('金門縣', '離島', 140000, 152),
    ('連江縣', '離島', 13000, 28)
) AS t(county, region, population, area_km2)

這張維度表用 VALUES 直接寫死 22 縣市的人口、面積、區域。資料來源是內政部戶政司與行政院主計總處的公開統計(讀者可依需要自行驗證),授權同樣為政府資料開放授權條款第 1 版。把這類「變動頻率極低」的維度資料寫死,可以避免每跑一次管線就拉一次外部 API。

-- de-journey/models/aqi_project/models/marts/mart_daily_summary.sql
{{ config(materialized='table') }}

-- 顆粒度 = 測站 × 日(給月報用)
SELECT
    DATE(observed_at)        AS observed_date,
    station_id,
    county,
    COUNT(*)                 AS reading_count,
    AVG(aqi)                 AS avg_aqi,
    MAX(aqi)                 AS max_aqi,
    SUM(CASE WHEN aqi > 100 THEN 1 ELSE 0 END) AS bad_hours,
    AVG(pm25)                AS avg_pm25
FROM {{ ref('fct_aqi_hourly') }}
WHERE observed_at IS NOT NULL
GROUP BY DATE(observed_at), station_id, county

這是聚合寬表,給月報與儀表板用。bad_hours 是「AQI > 100 的小時數」,搭配 reading_count 就能算出「不良日比例」(bad_hours / reading_count),這就是 Day 41 定義的 INDICATORS["bad_day_ratio"] 指標。把這個計算從 SQL 預先聚合好,儀表板讀 mart_daily_summary 就不用每次重新掃事實表。

把以上四個模型寫完後,跑一次 dbt run 確認全部成功:

cd de-journey/models/aqi_project
uv run dbt run

輸出(節錄):

1 of 7 OK  created stg_aqi_raw  ................... [OK in 0.42s]
2 of 7 OK  created stg_stations  .................. [OK in 0.31s]
3 of 7 OK  created int_aqi_with_station  .......... [OK in 0.18s]
4 of 7 OK  created fct_aqi_hourly  ............... [OK in 0.45s]
5 of 7 OK  created dim_stations  ................. [OK in 0.22s]
6 of 7 OK  created dim_counties  ................. [OK in 0.14s]
7 of 7 OK  created mart_daily_summary  ........... [OK in 0.38s]

七個模型跑完。fct_aqi_hourly 是事實表,所以 materialization 是 table(會在 DuckDB 落地成實體表);其他維度表也是 table;staging 與 intermediate 仍是 view。整個 dbt run 在本機 2 秒內跑完,效能遠超任何 OLAP 引擎。

第三步:把合約檢查接進 daily_run

Day 38 寫的合約檢查腳本已經能單獨跑了。今天把它接進 daily_run.py,讓 ingest、dbt run、合約檢查三件事變成一個完整的 daily_run 流程。我們把整個流程寫成幾個 Python 函式,方便排程器或 CI 呼叫。

"""de-journey/pipelines/daily_run.py(續):執行 ingest -> dbt run -> 合約檢查。"""
from __future__ import annotations

import duckdb
import subprocess

WAREHOUSE = Path(__file__).resolve().parents[1] / "warehouse" / "de-journey.duckdb"

def run_ingest(target_date: str) -> int:
    """呼叫 ingest 腳本並回傳 ingest 列數。"""
    cmd = ["uv", "run", "python", "pipelines/ingest_aqi.py",
           "--date", target_date, "--use-synthetic"]
    out = subprocess.run(cmd, capture_output=True, text=True, check=True)
    log.info("ingest 完成:%s", out.stdout.strip())
    n = int(out.stdout.strip().split()[-2])
    return n

def run_dbt() -> None:
    """在 dbt 專案目錄跑 dbt run,跑全部 7 個模型。"""
    dbt_dir = Path(__file__).resolve().parents[1] / "models" / "aqi_project"
    cmd = ["uv", "run", "dbt", "run"]
    subprocess.run(cmd, cwd=str(dbt_dir), capture_output=True, text=True, check=True)
    log.info("dbt run 完成")

def run_contract_check() -> dict:
    """跑 Day 38 的合約檢查,回傳 SLI 摘要。"""
    cmd = ["uv", "run", "python", "scripts/check_contract.py"]
    out = subprocess.run(cmd, capture_output=True, text=True, check=True)
    log.info("合約檢查完成")
    return {"status": "ok", "stdout_tail": out.stdout.splitlines()[-3:]}

def main(target_date: str) -> int:
    log.info("=== daily_run 開始:%s ===", target_date)
    n = run_ingest(target_date)
    run_dbt()
    sli = run_contract_check()
    con = duckdb.connect(str(WAREHOUSE))
    rows_aqi = con.execute(
        "SELECT COUNT(*) FROM aqi.fct_aqi_hourly"
    ).fetchone()[0]
    con.close()
    log.info("=== daily_run 結束:%d 列 ingest、%d 列 fct ===", n, rows_aqi)
    return 0

if __name__ == "__main__":
    p = argparse.ArgumentParser()
    p.add_argument("--date", default=datetime.now(timezone(timedelta(hours=8)))
                   .date().isoformat())
    args = p.parse_args()
    sys.exit(main(args.date))

這段把 daily_run 寫完。run_ingest 用 subprocess.run 呼叫 ingest 腳本並抓回 exit code 與輸出;run_dbt 在 dbt 專案目錄裡跑;run_contract_check 跑 Day 38 的腳本。每個步驟都有 try/except 與 logging,失敗時會留下明確的 log 訊息。輸出範例:

$ uv run python pipelines/daily_run.py --date 2025-12-09
=== daily_run 開始:2025-12-09 ===
ingest 完成:已 ingest 240 列(合成)
dbt run 完成
合約檢查完成
=== daily_run 結束:240 列 ingest、240 列 fct ===

這條管線現在已經能完整跑完一天的工作:240 列 ingest 進來、dbt 跑出 7 個模型、合約檢查通過、240 列事實表落地。整個流程在 CPU 上約 3 分鐘跑完(ingest 30 秒、dbt run 5 秒、合約檢查 1 秒,其餘是 Python 啟動與 DuckDB 連線)。

第四步:APScheduler 本機排程

APScheduler 3.11 是 Python 社群最常用的本機排程器,純 Python 實作、不需要外部服務。我們用 BackgroundScheduler 在背景跑每日 03:00 的 daily_run:

"""de-journey/pipelines/scheduler.py:本機 APScheduler 排程。"""
from __future__ import annotations

import logging
import sys
from datetime import datetime, timedelta, timezone
from pathlib import Path

from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.cron import CronTrigger

sys.path.insert(0, str(Path(__file__).resolve().parent))
from project_config import SCHEDULE  # noqa: E402
from daily_run import main as run_daily  # noqa: E402

log = logging.getLogger("scheduler")

def scheduled_job():
    today = datetime.now(timezone(timedelta(hours=8))).date().isoformat()
    log.info("排程觸發:daily_run --date %s", today)
    try:
        run_daily(today)
    except Exception as exc:
        log.exception("daily_run 失敗:%s", exc)

sched = BackgroundScheduler(timezone="Asia/Taipei")
sched.add_job(
    scheduled_job,
    CronTrigger.from_crontab(SCHEDULE["daily_run_cron"]),
    id="daily_aqi",
    name="每日 AQI 管線",
    max_instances=1,
    coalesce=True,
)
sched.start()
log.info("排程啟動:%s", SCHEDULE["daily_run_cron"])
# 輸出:排程啟動:0 3 * * *

這段把 APScheduler 接到 daily_run。CronTrigger.from_crontab() 從 project_config.SCHEDULE["daily_run_cron"] 讀 cron 表達式("0 3 * * *" 表示每天 03:00),這樣要改排程時間時只改 project_config.py。max_instances=1 防止前一次還沒跑完就觸發下一次;coalesce=True 把多次遺漏的觸發合併成一次。

實務上 APScheduler 跑在背景有兩種方式:直接執行這個腳本(適合本機開發),或包成 systemd 服務(適合伺服器)。對大多數資料團隊來說,本機用第一種、伺服器用第二種:

# 本機開發:直接跑(用 Ctrl+C 停止)
uv run python pipelines/scheduler.py

# 伺服器:用 systemd 包成服務(範例 unit 檔)
sudo tee /etc/systemd/system/aqi-scheduler.service <<EOF
[Unit]
Description=AQI Daily Pipeline Scheduler
After=network.target

[Service]
User=app
WorkingDirectory=/opt/de-journey
ExecStart=/opt/de-journey/.venv/bin/python pipelines/scheduler.py
Restart=on-failure

[Install]
WantedBy=multi-user.target
EOF
sudo systemctl enable --now aqi-scheduler.service

systemd unit 檔是伺服器部署的標準做法。Restart=on-failure 表示如果排程器掛了會自動重啟,這是 24 小時運轉的關鍵。User=app 限定執行身分(不要用 root);WorkingDirectory 設定正確的專案目錄。

第五步:GitHub Actions 雲端排程

APScheduler 適合本機與自架伺服器,但很多團隊偏好「雲端原生」的 GitHub Actions。它不需要維護主機、隨專案走、有現成的 cron 表達式語法。我們在 .github/workflows/ 下新增一支 workflow 檔:

# de-journey/.github/workflows/daily_aqi.yml
name: Daily AQI Pipeline

on:
  schedule:
    - cron: '0 3 * * *'   # 每日 03:00 UTC(等同 11:00 UTC+8,這裡示範用)
  workflow_dispatch:        # 允許手動觸發

jobs:
  daily-run:
    runs-on: ubuntu-latest
    timeout-minutes: 30
    steps:
      - uses: actions/checkout@v4

      - name: 安裝 Python 與 uv
        uses: astral-sh/setup-uv@v3
        with:
          python-version: '3.13'

      - name: 同步相依
        run: uv sync

      - name: ingest(合成模式示範,真實部署請移除 --use-synthetic)
        run: uv run python pipelines/ingest_aqi.py --date $(date -u +%Y-%m-%d) --use-synthetic

      - name: dbt run
        working-directory: models/aqi_project
        run: uv run dbt run

      - name: 合約檢查
        run: uv run python scripts/check_contract.py

      - name: 上傳 log
        if: always()
        uses: actions/upload-artifact@v4
        with:
          name: daily-run-logs
          path: logs/

這份 workflow 把 daily_run 的四個步驟(ingest、dbt run、合約檢查、上傳 log)做成 GitHub Actions 的步驟。cron: '0 3 * * *' 用 UTC 時間;台灣時間是 UTC+8,所以 03:00 UTC = 11:00 UTC+8,這裡只是示範,實務上你應該選在台灣深夜流量低的時段,例如台灣時間 03:00 = UTC 19:00(前一晚)。workflow_dispatch 讓你能從 GitHub UI 手動觸發這個 workflow,方便除錯。

uv sync 會從 uv.lock 把環境建起來,整個 CI 環境與本機開發完全一致。working-directory: models/aqi_project 把 dbt run 的工作目錄切到 dbt 專案根目錄。actions/upload-artifact@v4 把 logs/ 上傳成 artifact,方便事後下載查看。

第六步:失敗處理與重試

管線一定會失敗。Day 22 教過的「重試、警報、斷點續跑」三招在這裡都會用到。我們把 daily_run.py 加上重試邏輯:

"""de-journey/pipelines/retry.py:對 daily_run 加重試邏輯。"""
from __future__ import annotations

import logging
import sys
import time
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parent))
from daily_run import main as run_daily  # noqa: E402

log = logging.getLogger("retry")

def run_with_retry(target_date: str, max_attempts: int = 3, delay_sec: int = 300):
    """失敗時自動重試,最多 3 次、每次間隔 5 分鐘。"""
    for attempt in range(1, max_attempts + 1):
        try:
            log.info("嘗試 %d/%d", attempt, max_attempts)
            run_daily(target_date)
            log.info("嘗試 %d 成功", attempt)
            return 0
        except Exception as exc:
            log.warning("嘗試 %d 失敗:%s", attempt, exc)
            if attempt < max_attempts:
                time.sleep(delay_sec)
    log.error("重試 %d 次仍失敗,請人工介入", max_attempts)
    return 1

這個 retry.py 把 daily_run 包起來,加上重試邏輯。3 次重試、每次間隔 5 分鐘,是 Day 22 教的「指數退避」的簡化版。實務上可以改成「1 分鐘、5 分鐘、15 分鐘」三段間隔,讓重試間隔隨長更更長。當三次都失敗,最後印出 ERROR 等級的 log 與回傳非零 exit code,這樣 APScheduler 或 GitHub Actions 的 workflow 會把工作標記為失敗。

常見錯誤與踩雷

第一個雷:管線不冪等,跑兩次結果不同。我們的設計刻意把冪等性放在 ingest 階段:CREATE OR REPLACE TABLE 會把當天的資料完全覆寫,所以跑兩次 --date 2025-12-09 結果完全一樣。如果你寫的 ingest 是 INSERT INTO 而不是 CREATE OR REPLACE,第二次跑會出現重複列,這是 Day 20「水位標記」沒做好的典型症狀。

第二個雷:排程器時區沒設定。APScheduler 預設用 UTC,當你寫 cron='0 3 * * *' 它會在 UTC 03:00 觸發,而不是台灣時間 03:00。請一律顯式設定 timezone="Asia/Taipei",並用 cron 表達式的台灣時間解讀。GitHub Actions 的 cron 也是 UTC,註解一定要明寫「cron 在 UTC,換算成 UTC+8」。

第三個雷:dbt 找不到 profile。APScheduler 在背景執行時,環境變數與 ~ 路徑可能與互動 shell 不同。請把 dbt profile 用絕對路徑寫進 ~/.dbt/profiles.yml,或改用 DBT_PROFILES_DIR 環境變數指定。

第四個雷:管線跑成功但 SLI 沒寫。如果 check_contract.py 失敗了,daily_run 應該把 breach 寫進 logs/contract_breaches.jsonl 並以非零 exit code 退出,而不是「假裝沒事」。Day 38 的設計就是這樣,但實務上很容易在包 retry 時不小心把錯誤吃掉,請保持 raise 行為。

第五個雷:忘記給 cron 觸發預留寬限時間。GitHub Actions 的 cron 並不保證準點觸發,特別是尖峰時段可能延遲 10–30 分鐘。Day 41 的合約 freshness_hours: 26 就是預留的寬限:即使 cron 延遲 2 小時,SLI 仍能達標。設定 SLO 時要把這個延遲算進去。

效能與實務提醒

整條管線在 CPU 上跑完 240 列約 3 分鐘,這個速度對每天跑一次的場景綽綽有餘。如果未來要擴充到「每小時跑一次」或「資料量變 10 倍」,有三個加速方向:第一,把 ingest 改成多執行緒平行抓不同測站;第二,把 dbt 的 staging 與 intermediate 改成 materialized table 而非 view;第三,用 dbt run --select +mart_daily_summary 只跑必要模型。

實務上有三個取捨值得記得。第一,APScheduler vs GitHub Actions:APScheduler 適合本機與自架伺服器、可即時觸發、不用維護 GitHub 帳號;GitHub Actions 適合雲端原生、隨程式碼走、不用管理主機,但排程時間不保證準點。我們兩個都用:APScheduler 做本機執行、GitHub Actions 做雲端備援。第二,synthetic vs 真實資料:範例用合成是為了可重現,但實務部署一定要切回真實,並把 --use-synthetic 旗標移除。第三,retry.py 的 delay:5 分鐘是經驗值,太短會把真正的問題壓住、太長會讓資料延遲太久;可視業務需求調整。

小結

今天把專案篇的「管線實作與排程」蓋起來。我們寫了 ingest 腳本(含合成與真實兩種模式)、補上四個 marts 模型、把 daily_run 寫成單一入口、用 APScheduler 與 GitHub Actions 兩種方式設定排程,並加上重試邏輯。重點觀念有三個:第一,daily_run.py 是單一入口,無論手動、排程、CI 都呼叫同一支腳本;第二,project_config.py 是共用設定的單一來源,cron 表達式、模型清單、指標定義都集中管理;第三,冪等性是管線設計的核心,CREATE OR REPLACE TABLE 與合成資料的 seed 都服務這個原則。

明天,我們會進入專案篇的第三天「評估與迭代」。我們會把 Day 38 的合約檢查擴充成 dbt tests、自訂品質指標、A/B 評估,並用 DuckDB 跑回歸測試。整個評估流程也會沿用 project_config.INDICATORS 與 project_config.MODELS,確保 Day 41-45 的指標定義一致。

結語

今天的重點是「讓管線真的每天跑起來」。我們從 ingest 開始,串到 dbt run、合約檢查、排程與重試,把一條完整的每日管線做出來。讀完這篇你應該能回答:daily_run 的「成功完成」怎麼定義?APScheduler 與 GitHub Actions 怎麼分工?管線失敗時怎麼自動重試?明天,我們會把這套管線的「評估與迭代」機制做出來:怎麼判斷今天的資料品質夠好?怎麼從歷史指標找到瓶頸?

延伸資源

  • APScheduler 官方文件(3.11):https://apscheduler.readthedocs.io/en/3.x/。CronTrigger 與 BackgroundScheduler 的 API 與範例。
  • GitHub Actions 排程 cron 語法:https://docs.github.com/en/actions/using-workflows/events-that-trigger-workflows#schedule。注意 GitHub Actions cron 是 UTC,且延遲可能到 30 分鐘。
  • uv 官方安裝指南:https://docs.astral.sh/uv/。Day 2 已介紹,本篇用 uv sync 在 CI 重建環境。
  • Day 38 章節(監控與資料合約):今天的合約檢查是 Day 38 的延伸,把合約與 SLI 接到管線流程。
  • Day 22 章節(失敗處理):重試、警報與斷點續跑的設計原則,本篇的 retry.py 是其中一個實作。

留言

這個網誌中的熱門文章

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