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是其中一個實作。
留言
張貼留言