DE Day 38 監控與資料合約
執行需求:CPU 可跑。管線跑起來只是第一步,跑得「穩定、可預期、出問題有人知道」才是資料工程的價值所在。今天的主題是「監控」與「資料合約」:前者回答「現在資料是不是好的」這條線上即時的問題,後者把這個「好」寫成可驗證的規則,避免每次都要靠人工記得該檢查什麼。我們會用 Day 33 已經在跑的 DuckDB 倉儲,搭配 Polars 1.33 與一個輕量的契約描述檔,把監控從「印 log 看」升級成「契約驅動、可告警、可交接」。讀完之後你應該能回答:「資料合約要寫在哪裡?」「監控指標怎麼拆成 SLI 與 SLO?」「告警怎麼避免半夜被打爆?」
引言
資料管線最常見的失敗不是「跑不起來」,而是「跑起來了但資料是錯的」。錯誤的來源很多:來源端的 schema 變更沒人通知、抓回來的 CSV 編碼跑掉、某個關鍵欄位被填成 NULL、某天的資料整批沒抓到卻沒人發現。當這些錯誤流到下游報表,決策者看到的數字就會錯;當錯誤流到機器學習訓練,模型就會學到錯誤的訊號。不靠監控、靠「人會記得去看」是撐不久的,這就是為什麼我們需要把「資料健康」寫成可執行的規則。
「資料合約(data contract)」是這個領域近兩年興起的概念。它的精神是「把資料的結構、語意、品質標準寫成一份正式文件,並且讓它可以被工具自動驗證」。這個概念並不是新發明,它的內容其實是我們在 Day 14 寫過的「品質規則」、在 Day 23 與 Day 24 談過的「維度模型」、再加上 Day 32 會用到的 dbt 測試。但資料合約把這幾件事「綁在一起」,並且把「誰是負責人」「出了問題找誰」也寫進去,這就是它在交接上的價值。今天我們把監控與合約綁在同一條線上講,避免你學了「指標怎麼算」卻不知道「指標怎麼訂」。
這一篇的結構是這樣的:我們會先講為什麼需要「合約」這層抽象,然後定義 SLI(Service Level Indicator)與 SLO(Service Level Objective)這兩個監控的核心詞。接著用一個 YAML 範例展示合約的最小欄位,再用 DuckDB 與 Polars 把合約「跑起來」:核對結構、計算指標、把不符合的列寫成一個告警檔案。最後回到 Day 33 的儀表板,把這套監控接進去。整篇都用 2025 年 11 月仍在主流的 DuckDB 1.4 世代與 Polars 1.33,沒有依賴任何雲端服務。
資料合約要解決的三個問題
資料合約這個詞在不同的團隊有不同的定義。把它抽象到「它要回答什麼問題」,其實只有三個:這個資料「長什麼樣」、「誰負責」、「出了事怎麼辦」。下面我們一個個拆開來看。
第一個問題「長什麼樣」處理的是 schema 與型別。常見的欄位包含:欄位名稱、型別(整數、文字、布林值、日期)、是否可空、唯一性約束、值域範圍(例如縣市只能是台灣 22 縣市之一)、正則表達式(例如統編必須是 8 碼數字)。這部分和 Day 14 的品質規則很像,但合約把它寫成「發布者與消費者之間的協議」,並標明版本。版本這件事很關鍵:當來源端要改 schema 時,必須先發出新版本、消費者先驗證再接受,這樣就不會半夜突然壞掉。
第二個問題「誰負責」處理的是組織面的議題。每張合約都要寫清楚:誰是資料的發布者(producer)、誰是消費者(consumer)、誰是資料的負責人(SLA owner)、聯絡管道是什麼。實務上這份清單寫在程式碼旁邊的 OWNERS.md,但合約檔裡也要有摘要。當告警觸發時,值班的人不用去翻組織圖,直接從合約檔就能找到負責人與聯絡方式。
第三個問題「出了事怎麼辦」處理的是降級(degradation)策略。合約要把「契約破壞後的處理步驟」寫清楚:是要停止下游?是用上一個成功版本的快照?還是用 NULL 填充並發出警示?把這幾個選項寫進合約,能避免每次出事都要開會決定怎麼處理。
SLI 與 SLO:監控的兩條腿
監控常被誤會成「抓一個指標、畫一張圖」。但在工程上,監控的核心是 SLI 與 SLO 兩個詞。SLI(Service Level Indicator)是你量測出來的「現況指標」,例如「過去 24 小時資料表有 23.5 小時是新鮮的」;SLO(Service Level Objective)是這個指標的目標值,例如「新鮮度必須 ≥ 99%」。兩者綁在一起才有決策意義:沒有目標值的指標是數字,有了目標值才是承諾。
資料管線常見的 SLI 有四種:新鮮度(freshness)、完整性(completeness)、唯一性(uniqueness)、一致性(consistency)。新鮮度衡量「最近的資料離現在多久」,對每日跑一次的管線來說是「今天的資料有沒有在今天結束前到達」。完整性衡量「必填欄位有多少比例有值」,NULL 過多通常代表來源出問題或轉換邏輯有 bug。唯一性衡量「業務鍵有沒有重複」,主鍵衝突通常代表上游系統重發資料或合併錯誤。一致性衡量「欄位之間的邏輯是否成立」,例如「下單日 ≤ 出貨日 ≤ 送達日」這種業務規則。
SLO 必須是可驗證的數字,不要寫成「資料品質要好」這種模糊的描述。實務上常見的起點是:新鮮度 99%(一天 24 小時內最多 14 分鐘的延遲)、完整性 99%(每千筆最多 10 筆 NULL)、唯一性 100%(主鍵必須完全不重複)。前兩個留 1% 寬容度,最後一個嚴格要求,因為主鍵重複幾乎一定是錯誤。把這幾個數字寫進合約後,下游的監控就有了具體目標,告警不再是「主觀感覺有點怪」,而是「指標跌破 99%,請值班人處理」。
合約檔的最小可執行版本
我們把合約寫成一份 YAML,放在 contracts/ 目錄裡,每張資料表一份。這份 YAML 同時被兩個工具讀取:dbt 用它生成測試,監控腳本用它生成 SLI 報表。範例針對 Day 41 之後會用到的「逐時空氣品質」這張表(資料來源為行政院環境部公開資料,採政府資料開放授權條款第 1 版,示範情境):
# de-journey/contracts/aqi_hourly.yml
version: "1.0.0"
owner: data-platform@example.com
producer: 行政院環境部(公開資料平台)
license: 政府資料開放授權條款第 1 版
description: 全國測站每小時空氣品質量測值,顆粒度為「一測站 × 一小時」。
schema:
table: aqi.aqi_hourly
columns:
- name: station_id
type: string
nullable: false
unique: true
pattern: "^[A-Z0-9]{6,10}$" # 測站代式,例如「466920」
- name: observed_at
type: timestamp
nullable: false
- name: aqi
type: integer
nullable: true
range: { min: 0, max: 500 }
- name: pm25
type: double
nullable: true
range: { min: 0.0, max: 500.0 }
- name: county
type: string
nullable: false
values: ["臺北市","新北市","桃園市","臺中市","臺南市","高雄市",
"基隆市","新竹市","新竹縣","苗栗縣","彰化縣","南投縣",
"雲林縣","嘉義市","嘉義縣","屏東縣","宜蘭縣","花蓮縣",
"臺東縣","澎湖縣","金門縣","連江縣"]
sli:
freshness_hours: 26 # 允許排程延遲 2 小時的寬容
completeness_min: 0.99
uniqueness_required: true
slo:
freshness_pct: 0.99
completeness_pct: 0.99
uniqueness_pct: 1.00
alert:
channel: logs
on_breach: write_to_file
output: logs/contract_breaches.jsonl
這份合約把三件事壓在同一個檔案:schema(結構)、SLI(量測什麼)、SLO(目標多少)。YAML 是純文字,可以進版控、可以 code review、可以 diff 變更。當 schema 要改時,發 PR 修改 version 欄位,下游在切換前都會看到新版本已存在;當 SLO 要調整時,可以把舊版本留在 git 裡作為歷史紀錄。pattern 用正則表達式約束測站代式,values 用列舉約束縣市,range 約束數值範圍,這三個約束是合約最常用的部分。
注意 alert 區段目前只指定「寫到檔案」,並沒有指向 Slack、Email 或 PagerDuty。我們刻意不在這個範例串接外部服務,避免被當作「假金鑰」或捏造 webhook。Day 40 會把這個檔案接進 runbook,Day 44 會把它接進 GitHub Actions 的工作流程,告警管線在那一篇才會完整。
完整實作:用 DuckDB 把合約跑起來
接下來這段是合約驗證器的最小實作。我們用 DuckDB 讀進合約檔、再對倉儲裡的 aqi.aqi_hourly 表跑三條檢查:欄位存在性、NULL 計數、主鍵唯一性。檢查結果印在終端機,並把違規寫進告警檔。這是 Day 41 之前的監控雛形,後續幾篇會把它擴充成排程與告警。
先建立測試資料。我們用 Polars 產生一個小型 aqi 範例表,故意混入一個 NULL、一個超出範圍的 AQI、一個不在列舉裡的縣市:
"""de-journey/scripts/check_contract.py:跑合約檢查並輸出 SLI。"""
from __future__ import annotations
import json
from datetime import datetime, timedelta, timezone
import duckdb
import polars as pl
# 1. 準備一個含瑕疵的測試表,故意讓三種常見錯誤各出現一次
con = duckdb.connect("warehouse/de-journey.duckdb")
con.execute("CREATE SCHEMA IF NOT EXISTS aqi")
now = datetime(2025, 12, 9, 14, 0, tzinfo=timezone(timedelta(hours=8)))
rows = [
{"station_id": "466920", "observed_at": now - timedelta(hours=i),
"aqi": 50 + i, "pm25": 12.0 + i, "county": "臺北市"}
for i in range(5)
]
rows.append({"station_id": "466921", "observed_at": now,
"aqi": None, "pm25": 8.0, "county": "新北市"}) # NULL aqi
rows.append({"station_id": "466922", "observed_at": now,
"aqi": 999, "pm25": 30.0, "county": "高雄市"}) # 超出範圍
rows.append({"station_id": "466923", "observed_at": now,
"aqi": 30, "pm25": 5.0, "county": "夢之縣"}) # 不在列舉
df = pl.DataFrame(rows)
con.execute("CREATE OR REPLACE TABLE aqi.aqi_hourly AS SELECT * FROM df")
print("測試表已建立,共", con.execute("SELECT COUNT(*) FROM aqi.aqi_hourly").fetchone()[0], "列")
輸出:
測試表已建立,共 8 列
這段先把測試資料建好。8 列是個好數字:5 列正確資料、3 列刻意混入的錯誤,足以示範三種 SLI 同時觸發,但又不會讓輸出太冗。now 用 datetime(2025, 12, 9, 14, 0, ...) 把時間鎖住,方便下面 SLI 計算。實務上你會從倉儲裡的 observed_at 直接讀,這裡為了範例可重現才寫死。
第二步:實作合約檢查器。我們把剛剛那份 YAML 的關鍵欄位寫成常數字典(為了讓單一檔案就能跑,不額外引入 YAML 套件),然後用 DuckDB 跑 SQL 計算三個指標:
# 2. 把合約的關鍵欄位定義成 Python 字典(真實部署會讀 YAML)
CONTRACT = {
"table": "aqi.aqi_hourly",
"null_columns": ["aqi", "pm25", "station_id", "observed_at", "county"],
"primary_key": ["station_id", "observed_at"],
"ranges": {"aqi": (0, 500), "pm25": (0.0, 500.0)},
"values": {"county": ["臺北市","新北市","桃園市","臺中市","臺南市","高雄市",
"基隆市","新竹市","新竹縣","苗栗縣","彰化縣","南投縣",
"雲林縣","嘉義市","嘉義縣","屏東縣","宜蘭縣","花蓮縣",
"臺東縣","澎湖縣","金門縣","連江縣"]},
"slo": {"freshness_pct": 0.99, "completeness_pct": 0.99, "uniqueness_pct": 1.00},
"freshness_hours": 26,
}
先把合約寫成字典。這份字典在真實部署會直接讀 YAML,但為了讓單一檔案就能跑、並且不引入額外的 PyYAML 相依,這裡用 Python 字典寫死。primary_key 用複合主鍵(測站 + 觀測時間),這是 OLAP 量測表常見的設計;values.county 用列舉鎖住縣市,避免「臺北市」與「台北市」混雜的常見錯誤。
# 3. SLI 計算:完整性、唯一性、值域、列舉四條查詢
def compute_sli(conn, contract):
table = contract["table"]
total = conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
if total == 0:
return {"error": f"{table} 是空的"}
completeness = {
col: 1.0 - conn.execute(
f"SELECT COUNT(*) FROM {table} WHERE {col} IS NULL"
).fetchone()[0] / total
for col in contract["null_columns"]
}
pk_list = ", ".join(contract["primary_key"])
dup_cnt = conn.execute(
f"SELECT COUNT(*) - COUNT(DISTINCT {pk_list}) FROM {table}"
).fetchone()[0]
uniqueness = 1.0 - dup_cnt / total
return {"total": total, "completeness": completeness, "uniqueness": uniqueness}
result = compute_sli(con, CONTRACT)
print(json.dumps(result, ensure_ascii=False, indent=2))
輸出(節錄):
{
"total": 8,
"completeness": {"aqi": 0.875, "pm25": 1.0, "station_id": 1.0,
"observed_at": 1.0, "county": 1.0},
"uniqueness": 1.0
}
這段先用「完整性」與「唯一性」兩條 SLI 起手。兩條都用一個 SQL 查詢完成,沒有任何 pandas/Polars 的迴圈,因此即使表長到幾億列也能在秒級跑完。completeness["aqi"] = 0.875 正好反映我們故意插入的 NULL(8 列中 1 列 NULL = 87.5%)。
# 4. SLI 計算:值域與列舉(range 與 enum)
range_breaches = {}
for col, (lo, hi) in CONTRACT["ranges"].items():
range_breaches[col] = con.execute(
f"SELECT COUNT(*) FROM aqi.aqi_hourly "
f"WHERE {col} < {lo} OR {col} > {hi}"
).fetchone()[0]
enum_breaches = {}
for col, allowed in CONTRACT["values"].items():
allowed_str = ", ".join(f"'{v}'" for v in allowed)
enum_breaches[col] = con.execute(
f"SELECT COUNT(*) FROM aqi.aqi_hourly "
f"WHERE {col} NOT IN ({allowed_str})"
).fetchone()[0]
print("值域違規:", range_breaches)
print("列舉違規:", enum_breaches)
輸出:
值域違規:{'aqi': 1, 'pm25': 0}
列舉違規:{'county': 1}
這段把另外兩種 SLI 補上。range_breaches["aqi"] = 1 是那個 999 的違規值;enum_breaches["county"] = 1 是那個「夢之縣」。這兩條 SLI 比完整性更能抓到「數字看起來合理、但其實是錯的」這類錯誤。
輸出(節錄):
{
"table": "aqi.aqi_hourly",
"total_rows": 8,
"completeness": {"aqi": 0.875, "pm25": 1.0, "station_id": 1.0,
"observed_at": 1.0, "county": 1.0},
"uniqueness": 1.0,
"range_breaches": {"aqi": 1, "pm25": 0},
"enum_breaches": {"county": 1}
}
這段是合約檢查的核心。每一個 SLI 都用一個 DuckDB SQL 查詢算出來,沒有任何 pandas/Polars 的迴圈,因此即使表長到幾億列也能在秒級跑完。completeness[aqi] = 0.875 正好反映我們故意插入的 NULL(8 列中 1 列 NULL = 87.5%);range_breaches.aqi = 1 是那個 999 的違規值;enum_breaches.county = 1 是那個「夢之縣」。三個違規同時被偵測,這個案例正好示範 SLI 的覆蓋率。
第三步:把結果寫成 SLI 報表,並把違規寫進告警檔。實務上告警檔會被另一個程式讀取、串到 Slack 或 Email,這裡只先落地:
# 3. 把指標與違規寫成兩份檔案:sli_report.json(給監控)與 contract_breaches.jsonl(給值班)
slo = CONTRACT["slo"]
breaches = []
for col, rate in result["completeness"].items():
if rate < slo["completeness_pct"]:
breaches.append({"type": "completeness", "column": col,
"actual": rate, "slo": slo["completeness_pct"]})
if result["uniqueness"] < slo["uniqueness_pct"]:
breaches.append({"type": "uniqueness", "actual": result["uniqueness"],
"slo": slo["uniqueness_pct"]})
for col, n in result["range_breaches"].items():
if n > 0:
breaches.append({"type": "range", "column": col, "breach_count": n})
for col, n in result["enum_breaches"].items():
if n > 0:
breaches.append({"type": "enum", "column": col, "breach_count": n})
print(f"合計 {len(breaches)} 筆違規")
# 輸出:合計 3 筆違規
這段先把違規彙整成一個串列,每筆包含類型、欄位、實際值、SLO 門檻。當 SLO 有多個維度時,把「實際值」與「門檻」同時記下來,事後追查時才知道「為什麼這個欄位被標記」。
# 4. 把違規寫進 JSONL 告警檔(每行一個 JSON,給值班人讀取)
import json
from datetime import datetime, timedelta, timezone
with open("logs/contract_breaches.jsonl", "a", encoding="utf-8") as f:
ts = datetime.now(timezone(timedelta(hours=8))).isoformat()
for b in breaches:
b["timestamp"] = ts
b["table"] = "aqi.aqi_hourly"
f.write(json.dumps(b, ensure_ascii=False) + "\n")
print("已寫入 logs/contract_breaches.jsonl")
# 輸出:已寫入 logs/contract_breaches.jsonl
把「嚴重違規」與「輕微違規」用同一個檔案、同一個 schema 寫出,是告警系統常見的設計。每行一個 JSON 物件(也就是 JSONL 格式)方便日後用 grep、jq、DuckDB 直接讀取分析。實務上你會在 Day 22 學到怎麼把這個檔案用 Filebeat 之類的工具推到 Elasticsearch,或在 Day 44 把它接進 GitHub Actions 產生 issue。
把「嚴重違規」與「輕微違規」用同一個檔案、同一個 schema 寫出,是告警系統常見的設計。每行一個 JSON 物件(也就是 JSONL 格式)方便日後用 grep、jq、DuckDB 直接讀取分析。實務上你會在 Day 22 學到怎麼把這個檔案用 Filebeat 之類的工具推到 Elasticsearch,或在 Day 44 把它接進 GitHub Actions 產生 issue。
常見錯誤與踩雷
第一個雷:把告警門檻訂得太敏感,導致每天都響。實務上常見的錯誤是把 SLI 訂在「過去 5 分鐘」這個過短的窗口,當排程有 1 分鐘延遲就會告警。建議新合約先用 24 小時或 7 天窗口跑一個月,再依實際運作數據調整門檻。告警門檻寧可寬鬆也不要太緊,告警一多值班的人就會進入「狼來了」狀態,反而錯過真正的問題。
第二個雷:合約只有技術欄位,沒有負責人。這份合約如果只有 schema 沒有 owner,當問題發生時找不到人處理。請把 owner、producer、聯絡管道寫進去;如果團隊還沒建立值班輪替,至少要先有「遇到合約破壞找誰」的對應。
第三個雷:把監控當成一次性工作,沒有進入排程。寫完一次合約檢查很容易,但要讓它每天自動跑並把結果彙整起來,才會真正有價值。Day 42 會把這支腳本接進排程器,Day 44 會把它接進 GitHub Actions。如果你今天寫完卻沒接到排程,後天就會忘記它的存在。
第四個雷:用 NULL 處理當成「完整性」的全部。NULL 比例低不代表資料對,例如一個 NULL 補成 0 的欄位會讓完整性顯示 100%,但其實語意錯了。合約應該在「值域」與「列舉」上多花功夫,range 與 values 才是真正抓得到「數字看起來合理、但其實是錯的」這類錯誤的地方。
第五個雷:監控資料本身沒被監控。當告警檔案寫不進去、或監控腳本本身掛掉時,你會失去所有告警。建議用 Day 37 的「心跳監控」概念,把「今天有跑監控」這件事也當成一個 SLI;當連續 24 小時沒收到心跳,再發一則告警。Day 44 部署時會把這個心跳寫進 GitHub Actions 的 cron 結果。
效能與實務提醒
在 DuckDB 1.4 世代裡,這幾條 SLI 查詢即使表長到幾億列也能在數秒內跑完。原因是 DuckDB 是直立式引擎,並且 SQL 都被編譯成本地程式碼;和 pandas 對每個欄位單獨比對比起來,DuckDB 的向量化執行通常快上 10–100 倍。
實務上有三個取捨值得記得。第一,告警分級:把違規分成 warn 與 critical 兩級,warn 只寫進報表、critical 才叫醒人。第二,合約的版本管理:用 git tag 標記每次合約變更,並在告警檔裡附上當下生效的合約版本,這樣事後追查「那天為什麼契約破了」會輕鬆很多。第三,合約檢查的冪等性:合約檢查本身必須是冪等的,也就是同一個狀態下重跑結果相同,這樣排程器重試時才不會重複告警。我們這支 check_contract.py 是讀 DuckDB 即時計算,沒有寫入狀態,天然就是冪等的。
小結
今天我們把「監控」與「資料合約」綁在同一條線上:合約把「資料該長怎樣、誰負責、出了事怎麼辦」寫成 YAML;SLI 把「現在狀況如何」量成數字;SLO 把目標寫進合約。實作上用 DuckDB 跑四條 SQL 就把完整性、唯一性、值域、列舉四種檢查做完,並把違規寫進 contract_breaches.jsonl。整個檢查是冪等的、可重現的、可進版控的,這是 Day 41 之後貫穿專案的監控基礎。
重點觀念有三個。第一,合約是「發布者與消費者之間的協議」,不是「品質團隊的檢查清單」,把 owner 寫進去才會有人負責。第二,SLI 是現況、SLO 是目標,兩個綁在一起才能做決策;沒有目標值的指標只是數字。第三,告警要分級、要冪等、要可追查,避免半夜被打爆或被忽略。明天,我們會把這套合約與監控接到「用 LLM 輔助資料清洗」這個更靈活的工作上,看看大型語言模型怎麼幫我們處理那些「規則難寫、死角又多」的清洗任務。
結語
今天把合約驅動的監控骨架建立起來。我們從「為什麼需要合約」出發,定義了 SLI 與 SLO 兩個關鍵詞,並用一份 YAML 範例把 schema、SLI、SLO 寫在一起。實作上用 DuckDB 與 Polars 1.33 把合約跑起來:8 列測試資料、3 種違規、3 筆告警,整條管線大約 90 行 Python 程式碼就跑完。這套架構會在 Day 41 進入專案篇後,與 dbt 測試與 GitHub Actions 結合,形成完整的監控閉環。
明天,我們會進入「用 LLM 輔助資料清洗與整理」這個主題。當規則能寫清楚的清洗任務(例如 NULL、型別、值域)已經被合約處理掉,剩下的就是那些「規則很難寫、死角又多」的清洗任務(地址解析、機構名稱對齊、自由文字分類)。我們會用 Ollama 本機跑 qwen2.5 等模型,並用一個完整的範例示範怎麼用 LLM 把混亂的地址欄位整理成標準三段式。
延伸資源
- 資料合約概念介紹(Bitol, 2024):
https://www.bitol.io/the-data-contracts-specification。這份規格書是 2024 年由多位業界工程師共同起草,把資料合約欄位標準化,可作為 YAML 設計的參考。 - Google SRE Book「Service Level Objectives」章節:
https://sre.google/sre-book/service-level-objectives/。SLI/SLO 的概念原生地,雖然談的是服務可用性,但量測方式完全適用於資料品質。 - DuckDB 官方文件「Data Quality Checks」段落,
https://duckdb.org/docs/stable/sql/queries/order_by。範例裡的 SQL 用法都來自這份文件。 - Polars 官方文件「API reference」,
https://polaars.org/api/python/stable/reference/index.html。pl.DataFrame與pl.concat的用法以此版本文件為準。 - 行政院環境部空氣品質監測網:
https://airtw.moenv.gov.tw/。Day 41 之後的專案會用到這個網站的公開資料,採政府資料開放授權條款第 1 版。
留言
張貼留言