DE Day 33 端到端管線(四):失敗通知與重跑
執行需求:CPU 可跑。今天是端到端管線系列的第四天。昨天我們把品質檢查寫進管線,並把違規結果留在 meta.quality_check_result。今天要把這份紀錄變成「失敗通知」:當 critical 違規發生時,自動把違規摘要組合成一份通知訊息、用「模擬發信」的方式寫到 logs/notifications/ 目錄(標示為模擬,不串接真實郵件服務),並提供「單一資料集重跑」的 CLI 讓發生錯誤時不必整條管線重來。本篇所有範例都在 CPU 上執行。讀完這篇,你會有一套「出問題時知道、救得回來」的故障處理流程。
引言
管線跑久了,難免會遇到「昨天的資料突然變少」、「某個欄位被政府平台改成新格式」、「網路瞬斷讓下載失敗」。這些狀況不會每天發生,但發生時如果沒有對應的處理流程,使用者打開儀表板看到空白才回報問題,影響範圍通常已經擴大到「昨天的報表都沒了」。
今天要建立的故障處理流程有三個層次:第一,偵測——昨天的 check_quality.py 已經能做這件事;第二,通知——當 critical 違規發生時,把違規摘要用「模擬發信」方式送出(不串接真實 SMTP / Slack / LINE,避免本系列示範帳號洩漏問題);第三,重跑——提供 rerun.py CLI 讓發生問題時可以只重跑特定資料集,不必整條管線重來。實務上我們會把這三層串起來:通知裡附上「建議重跑的指令」,讓收到通知的人可以直接複製貼上。
注意:今天的所有「發信」動作都是「模擬」——我們會把通知訊息寫成 HTML 檔並存到 logs/notifications/ 目錄,標示清楚這是模擬、不是真實通知。這是「負責任的示範」原則:當一個開源範例包含「發通知」的程式碼時,不應該真的去串 SMTP,避免讀者不小心把測試信件寄給真實的人。
失敗通知與重跑的核心觀念
失敗通知的設計重點是「內容精簡、可行動」。一份好的通知應該包含四件事:發生時間、哪個資料集、哪條規則違規、違規多少筆、建議下一步動作。把這四件事寫成一行 80 字以內的訊息,使用者打開手機就能看完整,並能立刻採取行動。如果通知寫成一大篇,可能的問題是「使用者根本不會看完」,警報就失去意義。
另一個重要觀念是「通知的去敏感化」。管線裡的通知可能包含「公司統一編號」、「資本額」、「公司名稱」等資料。當我們要把通知寄到 Slack 或 LINE 時,這些資料可能違反公司資訊保護政策。我們的設計是把「資料本身」留在倉儲裡、通知只包含「統計數字與違規規則名稱」。這樣即使通知被誤傳,也不會洩漏敏感資料。
重跑設計的重點是「可指定範圍」。當昨天的維度表出問題時,使用者只想要重跑「昨天的維度表」而不是「整條管線」。我們會提供 --datasets company_basic 之類的命令列參數,讓重跑範圍可以細到「單一資料集」。同時 --from-stage transform 讓使用者指定「從轉換層開始跑」(不必重新下載昨天的 CSV)。這個 CLI 設計與 Day 21 介紹的排程腳本是互補的:排程器負責「全自動的日常執行」,CLI 負責「人工介入的例外處理」。
共用設定:擴充 common.py 加上通知設定
先把 pipelines/common.py 加上通知與重跑需要的常數。注意:所有聯絡資訊都是「模擬」用的佔位文字,實務部署時請用環境變數注入:
"""de-journey/pipelines/common.py:Day 30-35 共用的管線常數。"""
from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[1]
DATA_DIR = PROJECT_ROOT / "data"
WAREHOUSE_DIR = PROJECT_ROOT / "warehouse"
LOGS_DIR = PROJECT_ROOT / "logs"
NOTIFICATIONS_DIR = LOGS_DIR / "notifications"
DUCKDB_PATH = WAREHOUSE_DIR / "de-journey.duckdb"
DATASETS = {
"company_basic": {
"title": "公司登記基本資料",
"source": "moea_basic",
"table": "raw.company_basic",
"staging_table": "staging.company_basic_clean",
"dim_table": "mart.dim_company",
"key_columns": ["uniform_no", "company_name"],
"partition_prefix": "company_basic",
},
"company_change": {
"title": "公司變更登記資料",
"source": "moea_change",
"table": "raw.company_change",
"staging_table": "staging.company_change_clean",
"fact_table": "mart.fact_company_change",
"key_columns": ["uniform_no", "change_date", "change_item"],
"partition_prefix": "company_change",
},
}
QUALITY_RULES = {
"company_basic": [
("uniform_no_not_null", "uniform_no IS NOT NULL", "critical", 0),
("uniform_no_length_8", "LENGTH(uniform_no) = 8", "critical", 0),
("company_name_not_null", "company_name IS NOT NULL", "critical", 0),
("company_name_min_len", "LENGTH(company_name) >= 2", "warning", 100),
("capital_amount_non_negative", "capital_amount >= 0", "critical", 0),
("establish_date_valid", "establish_date IS NOT NULL", "warning", 1000),
],
"company_change": [
("uniform_no_not_null", "uniform_no IS NOT NULL", "critical", 0),
("change_date_not_null", "change_date IS NOT NULL", "critical", 0),
("change_item_not_null", "change_item IS NOT NULL", "critical", 0),
],
}
# 通知模擬設定(不會真的寄信,只寫到 logs/notifications/)
NOTIFY_CONFIG = {
# 模擬收件人;實務上用環境變數或 secrets 管理
"recipients": ["data-eng@example.invalid"],
# 主旨樣板;{dataset} 與 {date} 會被替換
"subject_template": "[DE Pipeline] {dataset} 品質異常 {date}",
# 標示為模擬,避免讀者誤以為會真的寄出
"is_simulation": True,
}
MIN_DAILY_ROW_RATIO = 0.90
HTTP_TIMEOUT_SEC = 30
RETRY_ATTEMPTS = 3
RETRY_BACKOFF_SEC = 2.0
這份擴充重點是三件事:第一,新增 NOTIFICATIONS_DIR 與 NOTIFY_CONFIG,把通知相關的常數集中;第二,is_simulation = True 是關鍵標記——所有呼叫通知程式的地方都應該先檢查這個 flag,避免在測試環境忘記切換;第三,保留 Day 30–Day 32 的所有設定,確保管線向下相容。
完整實作:模擬發信的通知程式
接下來寫 pipelines/notify.py。這支腳本讀取昨天的 meta.quality_check_result、組合成 HTML 通知、寫到 logs/notifications/,並在主控台印出「已模擬寄出」的訊息。執行前不需要安裝新套件。
"""de-journey/pipelines/notify.py:模擬寄出品質違規通知(標示為模擬)。
從 meta.quality_check_result 讀取當天的 critical 違規,組成 HTML 通知,
寫到 logs/notifications/,並在主控台印出「已模擬寄出」。不會真的呼叫
任何 SMTP / Slack / LINE 服務,避免測試環境誤寄。
"""
import logging
import sys
from datetime import date
from pathlib import Path
import duckdb
from pipelines.common import (
DATASETS,
DUCKDB_PATH,
NOTIFICATIONS_DIR,
NOTIFY_CONFIG,
)
NOTIFICATIONS_DIR.mkdir(parents=True, exist_ok=True)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
handlers=[logging.StreamHandler(sys.stdout)],
)
log = logging.getLogger("notify")
def fetch_breaches(con: duckdb.DuckDBPyConnection) -> list[dict]:
"""讀取今天的 critical 違規清單。"""
rows = con.execute("""
SELECT
check_date,
dataset,
rule_name,
severity,
violation_count,
threshold
FROM meta.quality_check_result
WHERE check_date = current_date
AND severity = 'critical'
AND is_breached
ORDER BY dataset, rule_name
""").fetchall()
return [
{
"check_date": r[0],
"dataset": r[1],
"rule_name": r[2],
"severity": r[3],
"violation_count": r[4],
"threshold": r[5],
}
for r in rows
]
def render_html(breaches: list[dict], today: date) -> str:
"""把違規清單組合成簡潔的 HTML 通知。"""
if not breaches:
return f"<html><body><h1>無違規</h1><p>{today} 全部資料集通過品質檢查。</p></body></html>"
rows_html = "\n".join(
f"<tr><td>{b['dataset']}</td><td>{b['rule_name']}</td>"
f"<td>{b['violation_count']}</td><td>{b['threshold']}</td></tr>"
for b in breaches
)
rerun_hint = "<br>".join(
f"重跑 {b['dataset']}:<code>python -m pipelines.rerun --datasets {b['dataset']}</code>"
for b in breaches
)
return f"""<html>
<head><meta charset="utf-8"><title>DE Pipeline 品質違規</title></head>
<body>
<h1>[模擬] 品質違規通知 {today}</h1>
<p>共 {len(breaks := breaches)} 條 critical 違規,請參考下表:</p>
<table border="1" cellpadding="4">
<thead><tr><th>資料集</th><th>規則</th><th>違規筆數</th><th>門檻</th></tr></thead>
<tbody>
{rows_html}
</tbody>
</table>
<h2>建議動作</h2>
<p>{rerun_hint}</p>
<hr>
<p style="color:gray;">這是模擬通知,由 pipelines/notify.py 產生。
實務部署請改用 SMTP / Slack / LINE webhook。</p>
</body></html>"""
def send_simulation(html: str, subject: str, recipients: list[str],
today: date) -> Path:
"""模擬寄出:把 HTML 寫到 logs/notifications/,回傳檔案路徑。"""
fname = NOTIFICATIONS_DIR / f"{today:%Y-%m-%d}-{subject.replace(' ', '_')}.html"
fname.write_text(html, encoding="utf-8")
log.info("[模擬寄出] subject=%s to=%s -> %s",
subject, ",".join(recipients), fname)
return fname
def main() -> int:
if not NOTIFY_CONFIG.get("is_simulation", True):
log.warning("is_simulation=False,但本範例未實作真實寄送,請補上 SMTP 程式碼")
today = date.today()
con = duckdb.connect(str(DUCKDB_PATH))
try:
breaches = fetch_breaches(con)
html = render_html(breaches, today)
# 每個資料集的主旨獨立一份(這樣可以分開寄給不同負責人)
for dataset_name in {b["dataset"] for b in breaches}:
subject = NOTIFY_CONFIG["subject_template"].format(
dataset=dataset_name, date=today
)
send_simulation(
html,
subject,
NOTIFY_CONFIG["recipients"],
today,
)
if not breaches:
log.info("今天沒有 critical 違規,不寄送通知。")
con.close()
return 0
except Exception:
con.close()
raise
if __name__ == "__main__":
sys.exit(main())
這支腳本的核心邏輯是 fetch_breaches() 與 render_html()。fetch_breaches() 從昨天的 meta.quality_check_result 撈出所有 critical 且 is_breached = true 的違規。SQL 用 ORDER BY dataset, rule_name 確保違規清單是穩定排序(同一組違規兩次跑的 HTML 必須一致)。
render_html() 把違規清單組合成 HTML 表格,並加上「建議動作」區塊把重跑指令嵌入通知。這是「可行動通知」的核心:使用者不需要再去查 wiki,直接從通知裡複製指令就能救。HTML 結尾的灰字「這是模擬通知」是負責任的示範,避免讀者誤以為會真的寄出。
send_simulation() 寫 HTML 到 logs/notifications/,檔名包含日期與主旨。這讓我們日後可以用 ls logs/notifications/ 看過去的通知歷史,進行「通知頻率分析」(例如「這個月 critical 違規突然變多,是不是資料源改了格式?」)。
完整實作:單一資料集重跑 CLI
通知提供了「知道問題」的能力,但要「救得回來」還需要「重跑」。我們寫一支 CLI 工具 rerun.py:
"""de-journey/pipelines/rerun.py:重跑單一資料集的指定階段。"""
import argparse
import importlib
import logging
import sys
from pipelines.common import DATASETS, DUCKDB_PATH
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
)
log = logging.getLogger("rerun")
# 階段與腳本的對應表(避免硬寫字串)
STAGE_SCRIPTS = {
"ingest": "pipelines.ingest",
"transform": "pipelines.transform_basic", # 簡化版,僅支援 company_basic
"build_marts": "pipelines.build_marts",
"check_quality": "pipelines.check_quality",
}
def rerun_dataset(dataset: str, from_stage: str) -> int:
"""從指定階段開始重跑單一資料集。"""
if dataset not in DATASETS:
log.error("未知資料集:%s(可用:%s)", dataset, list(DATASETS))
return 1
if from_stage not in STAGE_SCRIPTS:
log.error("未知階段:%s(可用:%s)", from_stage, list(STAGE_SCRIPTS))
return 1
log.info("[%s] 從 %s 階段開始重跑", dataset, from_stage)
# 為了簡化,這裡直接呼叫對應腳本的 main()
# 實務上會需要把「單一資料集」參數傳進去
script_name = STAGE_SCRIPTS[from_stage]
log.info("執行:python -m %s", script_name)
mod = importlib.import_module(script_name)
return mod.main()
def main() -> int:
parser = argparse.ArgumentParser(description="重跑管線的單一資料集")
parser.add_argument(
"--datasets",
required=True,
help="要重跑的資料集名稱(逗號分隔),例如 company_basic",
)
parser.add_argument(
"--from-stage",
default="ingest",
choices=list(STAGE_SCRIPTS),
help="從哪個階段開始重跑,預設 ingest",
)
args = parser.parse_args()
datasets = [d.strip() for d in args.datasets.split(",")]
overall_rc = 0
for ds in datasets:
rc = rerun_dataset(ds, args.from_stage)
overall_rc = overall_rc or rc
return overall_rc
if __name__ == "__main__":
sys.exit(main())
這支 CLI 接受兩個參數:--datasets(要重跑的資料集名稱,逗號分隔)與 --from-stage(從哪個階段開始,預設 ingest)。例如「從轉換層開始重跑 company_basic」可以這樣用:
python -m pipelines.rerun --datasets company_basic --from-stage transform
實務上每個階段腳本(ingest.py、transform_*.py、build_marts.py)都需要支援「只跑特定資料集」的參數。上面的範例用 STAGE_SCRIPTS 字典把階段對應到腳本,並透過 importlib 動態載入與呼叫 main()。這是「CLI 與排程共用同一套邏輯」的標準做法:CLI 與排程器都呼叫同一支腳本的 main(),避免「手動跑會通、排程跑會壞」的分歧。
回傳的 exit code 會被 shell 的 && 與 || 串接,實務上可以這樣用:
python -m pipelines.rerun --datasets company_basic --from-stage ingest \
&& python -m pipelines.notify
# 重跑成功才寄通知;失敗就停在重跑這步
完整實作:把通知與重跑串成一支「故障回應腳本」
把「重跑 + 驗證 + 通知」串成一支腳本 incident_response.py,讓收到通知的人可以直接呼叫:
"""de-journey/pipelines/incident_response.py:故障回應一鍵流程。
對單一資料集執行:
1. 從 ingest 階段重跑
2. 跑品質檢查
3. 如果還有違規,模擬寄出通知
執行:
python -m pipelines.incident_response --datasets company_basic
"""
import argparse
import importlib
import logging
import sys
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
)
log = logging.getLogger("incident")
def run_step(name: str, module_name: str) -> int:
log.info("[%s] 開始", name)
mod = importlib.import_module(module_name)
rc = mod.main()
if rc != 0:
log.error("[%s] 失敗,停止後續步驟", name)
return rc
def main() -> int:
parser = argparse.ArgumentParser()
parser.add_argument("--datasets", required=True)
parser.add_argument("--from-stage", default="ingest")
args = parser.parse_args()
# 第一步:重跑(從指定階段開始)
rc = run_step(
f"rerun {args.datasets}",
"pipelines.rerun",
)
if rc != 0:
return rc
# 第二步:品質檢查
rc = run_step("check_quality", "pipelines.check_quality")
if rc != 0:
# 還是有違規 → 通知
log.info("品質檢查仍有違規,呼叫通知程式")
run_step("notify", "pipelines.notify")
log.info("故障回應流程完成")
return 0
if __name__ == "__main__":
sys.exit(main())
這支腳本展示「流程編排」的概念:把多支腳本按順序串起來,每步失敗就停止。實務上 Day 36 與 Day 37 會用 Airflow / GitHub Actions 把這套流程自動化;今天先用 Python 的 importlib 做到「命令列版的流程編排」,讓收到通知的人可以一鍵處理。
如果要進一步把「通知模板」抽成可設定的 YAML,方便不同環境調整內容:
"""de-journey/pipelines/notify_yaml.py:把通知模板抽到 YAML。"""
from pathlib import Path
import yaml
# YAML 檔案位置:pipelines/templates/notify.yaml
template_path = Path(__file__).parent / "templates" / "notify.yaml"
config = yaml.safe_load(template_path.read_text(encoding="utf-8"))
# 設定檔裡可以包含主旨、收件人、HTML 樣板
print("主旨樣板:", config["subject_template"])
print("收件人:", config["recipients"])
print("HTML 樣板長度:", len(config["body_template"]))
這支小腳本展示「設定與程式碼分離」的延伸做法:把主旨、收件人、HTML 樣板抽到 YAML,不同環境(測試 vs 正式)只需要替換 YAML 檔。實務上建議用 PyYAML 解析;如果還沒裝,uv pip install pyyaml 即可。
"""de-journey/pipelines/parse_notification_log.py:解析歷史通知,做頻率分析。"""
import re
from collections import Counter
from pathlib import Path
from pipelines.common import NOTIFICATIONS_DIR
# 用 regex 從 HTML 標題抓出日期
pattern = re.compile(r"品質違規通知 (\d{4}-\d{2}-\d{2})")
counter = Counter()
for html_file in NOTIFICATIONS_DIR.glob("*.html"):
m = pattern.search(html_file.read_text(encoding="utf-8"))
if m:
date = m.group(1)
counter[date] += 1
for date, n in sorted(counter.items())[-7:]:
print(f"{date}: {n} 份通知")
這支小腳本展示「通知歷史」的可分析性:當 logs/notifications/ 累積一段時間後,可以用 regex 從 HTML 標題抓出 (日期, 通知數) 的時間序列,進而判斷「這個月違規是不是變多了」。如果想做得更完整,可以把這個統計結果回寫到 meta.notification_count_daily 表,搭配 Day 32 的品質趨勢做交叉分析。
常見錯誤與踩雷
錯誤一:忘了設 is_simulation = True,導致在測試環境把測試信件寄給真實的人。常見症狀:某位同事收到一封主旨奇怪的郵件,內含「這是模擬通知」字樣。對應排查方向:所有通知相關的設定都應該預設 is_simulation = True,只有在正式部署時用環境變數切換;同時 send_simulation() 的函式命名也要明確標示「simulation」,避免有人誤把它當成真實寄送。
錯誤二:HTML 通知太長,使用者看完才發現沒有「建議動作」。常見症狀:使用者反映「警報太多看不懂」。對應排查方向:通知的「建議動作」必須放在 HTML 的開頭(一進來就看到),而不是埋在表格後面。我們用 標題讓這一區塊視覺上突出。建議動作
錯誤三:--from-stage 參數沒有對應的階段腳本,STAGE_SCRIPTS 字典少寫一筆。常見症狀:KeyError: 'transform_change'。對應排查方向:用 argparse 的 choices 參數限制合法值;並在 STAGE_SCRIPTS 字典加上新階段時同步更新測試。
錯誤四:importlib.import_module() 找不到模組。常見症狀:ModuleNotFoundError: No module named 'pipelines.ingest'。對應排查方向:importlib.import_module() 需要 pipelines 套件路徑在 sys.path 裡;執行時一定要用 python -m pipelines.rerun 形式呼叫,而不是 python pipelines/rerun.py。
錯誤五:重跑之後忘記跑品質檢查就標記「已修復」。常見症狀:通知說「已重跑」,但第二天又收到同樣的違規。對應排查方向:incident_response.py 的第二步永遠是 check_quality;不要為了「通知趕快關掉」而跳過驗證。
效能與實務提醒
通知與重跑的效能瓶頸在「重跑階段」。如果從 ingest 開始重跑,整條管線跑完約 3–5 分鐘(包含下載、DuckDB 寫入、轉換、品質檢查);如果從 transform 開始,則只要 30 秒。實務上我們會建議「先試 --from-stage transform,不行才回到 --from-stage ingest」的策略:90% 的問題是轉換層而非落地層,重跑轉換省下大量時間。
另一個工程建議:把 incident_response.py 的 exit code 當成「KPI」:每週統計「這個月有幾次故障回應」、「平均修復時間」是多少。當你看到「平均修復時間從 15 分鐘降到 5 分鐘」,代表 CLI 與通知流程確實有效;如果平均時間越來越長,代表資料源越來越不穩定。
實務上還有一個重要的取捨:「通知頻率」。當 critical 違規每天都有,使用者會疲乏,最後忽略所有通知。我們建議的策略是:把「每天都會發生的輕微違規」降級為 warning,只寄「異常變化」的 critical(例如「今天突然多了 10 萬筆違規」才寄)。這個動態門檻的設計會在 Day 38 的「資料合約」章節進一步討論。
小結
今天為端到端管線加上了失敗通知與重跑。我們擴充了 pipelines/common.py 的 NOTIFY_CONFIG;用 notify.py 把 critical 違規組合成 HTML 通知並模擬寫到 logs/notifications/;用 rerun.py CLI 提供「單一資料集、從指定階段開始」的重跑能力;最後用 incident_response.py 把「重跑 → 驗證 → 通知」串成一鍵流程。重點回顧:第一,通知內容要精簡、可行動,並放在 HTML 開頭;第二,is_simulation = True 是模擬發信的關鍵標記,避免測試環境誤寄;第三,importlib 動態載入是「CLI 與排程共用邏輯」的標準做法;第四,--from-stage 參數讓重跑範圍可控;第五,故障回應的 KPI(平均修復時間)應該每週追蹤。
明天 Day 34 會接著做「Streamlit 儀表板」:用 mart.dim_company 與 mart.fact_company_change 當資料源,建立一個可以讓使用者「按統一編號查公司」、「按月份看變更趨勢」、「看品質違規的時間序列」的互動介面。Streamlit 1.5x 是 2025 年主流的 Python 儀表板框架,本篇會用最少的程式碼完成一個可用的儀表板雛形。
結語
今天的重點是「讓管線從跑得對變成出問題時救得回來」。我們沒有引入 SMTP / Slack / LINE 套件(避免真實服務帳號洩漏),而是用「模擬寫到檔案」的方式展示通知的完整流程。這是負責任的開源示範原則:當一個範例包含「發通知」的程式碼時,必須明確標示「這是模擬」,並提供「真實部署時怎麼改」的指引。
明天,我們會把 mart 兩張表接到 Streamlit 1.5x,建立一個可以查公司、看趨勢、查品質結果的儀表板。Streamlit 的 API 很直觀:每行程式碼都是一段 UI 行為;改一行就能看到畫面更新。我們會用一個可以整段執行的範例,從「打開首頁看到 KPI 卡片」到「按統一編號查到單一公司的所有變更」一步步帶過。
延伸資源
- SMTP 官方文件(2025):
https://docs.python.org/3/library/smtplib.html。本篇範例刻意不串接真實 SMTP;實務部署時請用smtplib.SMTP()+ 環境變數管理帳號密碼,並加上 TLS。 - Slack Incoming Webhook(2025):
https://api.slack.com/messaging/webhooks。另一個常見的通知管道,requests.post()就能呼叫;同樣需要從環境變數讀 webhook URL。 - LINE Notify 官方文件(2025):
https://notify-bot.line.me/doc/en/。台灣團隊常用 LINE 當通知管道;token 一樣從環境變數讀。 - argparse 官方教學(2025):
https://docs.python.org/3/howto/argparse.html。本篇rerun.py的命令列設計以此文件為準。choices參數能自動生成 help 訊息並限制合法值。 - importlib 官方文件(2025):
https://docs.python.org/3/library/importlib.html。本篇用importlib.import_module()動態載入階段腳本,避免硬寫 import 字串。 - 政府資料開放平臺公司登記資料檢視頁:
https://data.gov.tw/。本系列使用的「公司登記資料」與「公司變更登記資料」位於此平台,授權為「政府資料開放授權條款第 1 版」。
留言
張貼留言