AI 互動教室 ‹ MLOps 自動化技術
下載 .py 開啟實戰 notebook ↗ 留言回報
CAPSTONE · DAGSTER × MLFLOW · 05

Dagster × MLflow:
一條會自己訓練、評估、把關、上線的管線

前四課你各拿到一個零件:訓練有紀錄、最好的版本上得了線、資料查得到來歷、有人會按執行。 這一課把它們焊成一條線,並且加上最重要的那個零件——一道會說「不」的閘門。 下面是這條線在 notebook 裡真的跑過的四次,按一顆按鈕看它走一遍:

四次執行、同一份程式碼,改的只有設定

每一個數字都是 notebook 的實測結果(同一組亂數種子):AUC 0.9684 → 0.9508 → 0.9698 → 0.8641, 其中兩次被閘門擋下、沒有任何東西上線。

01 · RESOURCE

先把「MLflow 在哪」變成可以注入的設定

class MlflowResource(dg.ConfigurableResource): tracking_uri: str experiment: str def setup(self) -> None: mlflow.set_tracking_uri(self.tracking_uri) # 帳本記在哪 if mlflow.get_experiment_by_name(self.experiment) is None: mlflow.create_experiment(self.experiment, artifact_location=...) mlflow.set_experiment(self.experiment) # run 歸到哪個實驗 @dg.asset def trained_model(context, config: TrainConfig, mlflow_res: MlflowResource, train_test: dict) -> str: mlflow_res.setup() # ← 這一行決定這次訓練被記到哪裡 ...

第 4 課的 resource 在這裡派上真正的用場。MLflow 的 tracking server 是典型的「環境」: 筆電上是一個 SQLite 檔、正式環境是內網的伺服器、CI 上又是另一個。 把它宣告成 ConfigurableResource(一個 Pydantic 模型,欄位型別會被檢查), 資產只要在參數列寫 mlflow_res: MlflowResource 就會被注入—— 換環境只改一行 Definitions,用得到 MLflow 的那幾個資產與檢查,程式一個字都不用動。

忘了 setup() 會怎樣?實測的答案很嚇人:Dagster 回報 success = True、沒有任何錯誤,但管線的 tracking 裡一個 run 都沒有, 工作目錄多出一個誰也不會去看的 mlflow.db——訓練確實跑了,只是記到了別的地方。 這種「安靜地做錯」正是 MLOps 最貴的一類 bug。

到 notebook 的 1️⃣ 節:把 MLflow 做成 resource
02 · 資產鏈

訓練資產:兩邊互相記下對方的 id

with mlflow.start_run(run_name=f"dagster-{context.run_id[:8]}") as run: mlflow.log_params({"model": ..., "max_depth": ..., "drift": ..., "dagster_run": context.run_id}) # MLflow 記住 Dagster clf.fit(X, y) info = mlflow.sklearn.log_model(clf, name="churn_model", signature=infer_signature(X, clf.predict_proba(X)[:, 1]), input_example=X.head(3)) mlflow.set_tag("dagster.asset", "trained_model") context.add_output_metadata({"mlflow_run": run.info.run_id, # Dagster 記住 MLflow "model_uri": info.model_uri}) return info.model_uri # 資產的內容是地址,不是模型物件

這一格是整條管線的縫合處:第 1 課的 start_run 疊在第 3 課的 @asset 上。真正的關鍵是雙向——只記單向的話,總有一天你會站在錯的那一邊: 有人看到線上模型怪怪的,是從 Registry 往回查;資料工程師發現某天的資料有問題,是從 Dagster 往前查。

資產回傳的是一個字串 models:/m-…,不是模型物件。模型檔案已經在 MLflow 那裡了, 管線裡傳的是它的地址;下游要用就自己載——這也讓「資產」在真實部署裡不必扛著幾百 MB 的東西跑。

到 notebook 的 2️⃣–3️⃣ 節:資料資產與訓練資產
03 · 評估

一行 evaluate,指標同時進兩本帳

with mlflow.start_run(run_name="evaluate"): res = mlflow.models.evaluate(trained_model, test_df, targets="label", model_type="classifier") metrics = {k: float(v) for k, v in res.metrics.items() if k in ("roc_auc", "accuracy_score", "f1_score", "recall_score", "precision_score")} context.add_output_metadata({k: dg.MetadataValue.float(v) for k, v in metrics.items()}) return metrics

第 2 課的 evaluate 一行產出 8 個指標與 5 張圖,全部記進當前 run。這裡多做一件事: 把其中幾個掛到 Dagster 的中繼資料上。重複記不是浪費——兩本帳的讀者不同: MLflow 的指標是拿來比較實驗的(20 個 run 排序找最好的那個), Dagster 的中繼資料是拿來看管線的(在資產圖上一眼看到這份指標上次算出來多少,還會畫出歷次趨勢)。

這個資產同時吃 trained_model(一個 URI 字串)與 train_test(切好的資料), 參數名各自對到上游——評估用的是 test 那一半,訓練資產從來沒看過它。

到 notebook 的 4️⃣ 節:評估資產
04 · 品質閘

為什麼閘門是 asset check,不是一個 if

@dg.asset_check(asset=model_metrics, blocking=True, description="品質閘:AUC 必須 ≥ 0.95 而且不輸目前 champion") def quality_gate(mlflow_res: MlflowResource, model_metrics: dict) -> dg.AssetCheckResult: mlflow_res.setup() try: champ = MlflowClient().get_model_version_by_alias("churn-clf", "champion") champion_auc = mlflow.get_run(champ.run_id).data.metrics.get("eval_auc", 0.0) except MlflowException: # 第一次執行:Registered Model ... not found champion_auc = 0.0 auc = model_metrics["roc_auc"] return dg.AssetCheckResult(passed=bool(auc >= 0.95 and auc >= champion_auc), severity=dg.AssetCheckSeverity.ERROR, metadata={"auc": auc, "champion_auc": champion_auc, "min_auc": 0.95})

兩道條件缺一不可:絕對門檻(不管以前多爛,低於 0.95 就是不准上線) +相對門檻(不能比現在線上那版還差)。你當然可以把這段寫成上線資產開頭的一個 if,但那樣會失去四件事:

寫成 if寫成 blocking asset check
這次過了沒埋在日誌裡,要翻一筆檢查結果進帳本,紅叉綠勾
為什麼沒過要自己 printmetadata 帶著 auc/champion_auc/門檻,永久保存
下游會不會跑每個下游都要再判一次blocking 直接擋住全部下游
這次執行算成功嗎算成功(你自己 return 了)run 標記為失敗,該叫的人會被叫

最後一項最關鍵:閘門擋下來時,這次執行必須是失敗的。如果它算成功,你的監控就永遠不會響—— 一條「安靜地什麼都沒做」的管線,比一條會壞掉的管線危險得多。 另外記得 dg.materialize() 沒有 asset_checks= 參數: 檢查要跟資產放同一個清單,忘了放它就靜靜地不執行,而執行結果照樣顯示成功。

到 notebook 的 5️⃣ 節:品質閘與 blocking 的效果
05 · 上線

放行之後:註冊、補記、移動 alias

mv = mlflow.register_model(trained_model, MODEL_NAME) # 1. 註冊成新版本 src_run = client.get_model_version(MODEL_NAME, mv.version).run_id client.log_metric(src_run, "eval_auc", model_metrics["roc_auc"]) # 2. 補記評估分數 client.set_registered_model_alias(MODEL_NAME, "champion", mv.version) # 3. 移動 alias = 上線 client.set_model_version_tag(MODEL_NAME, mv.version, "dagster_run", context.run_id)

第 3 步之後,服務端那行 load_model("models:/churn-clf@champion") 一個字都不用改, 下次載入就是新版。這就是第 2 課 alias 的用途,只是現在按下它的不是你的手,是管線。

第 2 步是最容易漏掉的一步。評估是在另一個名叫 evaluate 的 run 裡做的, 而下一次執行時,閘門要問的是「現任 champion 當初考幾分」——它讀的是 champion 版本指向的訓練 run。 沒把 eval_auc 補記上去,閘門讀到的永遠是實測的 metrics = {} → .get('eval_auc', 0.0) = 0.0, 相對門檻形同虛設:任何 AUC ≥ 0.95 的模型都能把 champion 換掉,包含比現任差的那些。 不報錯、不噴紅字,只是品質閘默默失效。

到 notebook 的 6️⃣ 節:上線資產
06 · 四次執行

同一份程式碼,四種劇情

執行設定AUC當時 champion結果
run 1RandomForest depth 80.9684—(Registry 空的)通過 → 成為 v1
run 2LogisticRegression0.95080.9684過了 0.95,但輸給現任 → 擋
run 3RandomForest depth 160.96980.9684通過 → 晉升 v2
run 4depth 16 + drift 1.50.86410.9698低於絕對門檻 → 擋

被擋的兩次,registered_champion 從實體化清單裡整個消失,執行結果是失敗—— 這正是 blocking=True 在做的事。

最值得停下來想的是 run 4:它跟 run 3 的模型設定一模一樣,同樣的森林、同樣的深度、同樣的種子。 程式碼沒動、參數沒動,AUC 卻掉了 0.10。變的只有資料。真實世界最常見的模型事故就是這樣: 沒有人改壞任何東西,是世界變了。閘門的價值不在於它擋下的那些模型,而在於你不必再靠運氣。

到 notebook 的 7️⃣ 節:連跑四次,看閘門開關
07 · 追溯

「現在線上那個模型,是怎麼來的?」

Registry @champion ──▶ MLflow run ──▶ dagster_run 參數 ──▶ Dagster 那次執行 │ materialization 的 mlflow_run 中繼資料 ◀─────────────────┘
champ = client.get_model_version_by_alias("churn-clf", "champion") # v2 run = mlflow.get_run(champ.run_id) # 參數、指標、tag dagster_run_id = run.data.params["dagster_run"] # → Dagster 那次執行 instance.get_run_by_id(dagster_run_id).status # SUCCESS records = instance.fetch_materializations(dg.AssetKey("trained_model"), limit=10).records match = next(r for r in records if r.run_id == dagster_run_id) match.asset_materialization.metadata["mlflow_run"].value == champ.run_id # True

半夜有人問「線上跑的模型是誰、什麼時候、用什麼資料訓的」,這題就是三跳到底的查詢: 從 alias 找到版本、從版本找到 MLflow run、從 run 的參數找到 Dagster 的執行 id; 反過來也走得通——用那個執行 id 去 Dagster 的帳本裡撈當時的中繼資料,會拿到同一個 MLflow run。 實測兩個方向指到同一個 run:True

在正式環境裡這兩跳都是 UI 上的一個連結。自己用 API 走一遍的意義是:你知道那個連結底下是什麼, 也知道當初少記一邊的 id,今天就會斷在哪裡。

到 notebook 的 8️⃣ 節:兩個方向各查一次
08 · 誰來按執行

收尾:把這條線交給排程、感測器與自動化條件

train_job = dg.define_asset_job("nightly_train", selection=dg.AssetSelection.all()) nightly = dg.ScheduleDefinition(name="nightly_train_schedule", job=train_job, cron_schedule="0 3 * * *", execution_timezone="Asia/Taipei") @dg.sensor(name="new_data_sensor", job=train_job, minimum_interval_seconds=30) def new_data_sensor(context): # 收件匣有新檔案就重訓,cursor 記住看過幾個 seen = int(context.cursor) if context.cursor else 0 files = sorted(INBOX.glob("*.csv")) if len(files) > seen: context.update_cursor(str(len(files))) yield dg.RunRequest(run_key=f"batch-{len(files)}") else: yield dg.SkipReason(f"沒有新資料(已處理 {seen} 批)") @dg.asset(automation_condition=dg.AutomationCondition.eager()) # 上游一更新,我自己就該重算 def data_profile(churn_data: pd.DataFrame) -> dict: ... # 有新資料就重出資料剖析 @dg.asset(automation_condition=dg.AutomationCondition.eager()) def champion_scorecard(registered_champion: str) -> str: ... # champion 換人就重出成績單 production_defs = dg.Definitions(assets=[...], asset_checks=[quality_gate], resources=RESOURCES, jobs=[train_job], schedules=[nightly], sensors=[new_data_sensor])

第 4 課的三種零件,回答的是同一個問題的三種版本:誰來按? 排程是「時間到了」、感測器是「有事情發生了」、自動化條件是「我的上游更新了,我自己該動了」。 在 notebook 裡用 evaluate_tick() 直接問它們「假設那個時刻到了,你會發什麼單子」—— 不必啟動任何背景服務就看得到答案:排程的 tick 送出 1 張 RunRequest; 感測器三個 tick 依序是「沒有新資料」→「發出 batch-1」→「不重複觸發」。

自動化條件的實測結果更值得看。在一本乾淨的帳本上跑一次會被閘門擋下的管線,然後連問三個 tick: 第一個 tick 是 0(評估器要先有基準,之後才知道什麼是新的——正式環境的 daemon 一直在跑、基準早就有了); 第二個 tick 只有 data_profile 舉手,因為它的上游 churn_data 剛更新過; 而 champion_scorecard 一動也不動——它的上游被閘門擋住、這次根本沒有產出。 閘門擋下的不只是這一次執行的下游,連自動化都跟著停在那裡,這是品質閘最容易被低估的一面。

最後全部收進一份 Definitions:7 個資產、1 個檢查、1 份資源、1 個 job、1 個排程、1 個感測器。 把它存成專案裡的 definitions.py,然後 dagster dev—— 瀏覽器打開就是資產圖、實體化歷史、檢查的紅叉綠勾、排程與感測器的開關。到這裡,這條線就不需要你了。

到 notebook 的 9️⃣ 節:job、排程、感測器、Definitions
09 · 實戰

換你動手

LEVEL 1

MIN_AUC 改成 0.97、換一個乾淨的模型名字,再跑一次 rf depth 8(AUC 約 0.968)。 它應該再也上不了線——確認 registered_champion 從實體化清單裡消失了。門檻是一個決策,不是一個常識。

LEVEL 2

再加一道 blocking 檢查 recall_gate(recall ≥ 0.9),跟品質閘掛在同一個資產上, 用 LogisticRegression(recall 約 0.86)跑一次看它被哪一道擋下來。順便把其中一個檢查從清單裡拿掉再跑一次, 體會「執行成功,但閘門根本沒跑」。

LEVEL 3

churn_data 改成 DailyPartitionsDefinition 的分割資產,讓每天只訓當天的資料, 再用 build_schedule_from_partitioned_job 產生排程。順便想一個沒有標準答案的問題: 分割之後,每一片各自要跟誰比?

卡住了?三題在 notebook 末節都有折疊解答與驗證方式——先自己做,再打開對照。

10 · 驗收

情境測驗

離開前試試看:下面的情境都真的會遇到。每題選一個你認為的最佳做法,選了馬上看得到解釋。

Q1 情境題

你要在自動重訓管線裡加一道「AUC 不到 0.95 就不准上線」的規則。以下哪種做法最好?

blocking 的資產檢查一次給你四件事:判斷結果存成一筆可查詢的紀錄(UI 紅叉綠勾)、判斷依據留在 metadata 裡(auc/champion_auc/門檻)、下游全部自動不跑、而且整次執行標記為失敗——監控才會響。A 的致命傷正是「執行不會失敗」:一條安靜地什麼都沒做的管線,比一條會壞掉的管線危險得多。B 把人放回迴圈裡,規模一大就變成瓶頸,而且沒有紀錄。D 雖然會失敗,但錯誤混在例外堆疊裡、沒有結構化的判斷依據,而且評估資產本身會被標成失敗——它其實算對了,是模型不合格。

Q2 錯誤診斷

管線跑了三個月都很正常,直到有人發現線上模型比上個月的還差。你去查閘門的中繼資料,看到每一次都長這樣。問題出在哪?

quality_gate passed=True {"auc": 0.9508, "champion_auc": 0.0, "min_auc": 0.95} # 而 champion 那個版本的訓練 run: mlflow.get_run(champ.run_id).data.metrics # {} → .get("eval_auc", 0.0) = 0.0

champion_auc 一直是 0.0 就是證據:閘門讀的是 champion 版本指向的訓練 run,而評估是在另一個 run(evaluate)裡做的。上線時要用 log_metriceval_auc 補記到訓練 run 上,下一次的閘門才比得到——漏了這一步,只剩絕對門檻,任何過 0.95 的模型都能把更好的 champion 換掉。A 症狀不符:alias 抓的版本是對的,是那個 run 上沒有指標;C 這裡的 try/except 是為了「第一次執行還沒有 champion」而存在的,拿掉只會讓第一次直接爆炸;D 只是把標準拉高,相對門檻依然失效——今天過 0.97 的爛模型照樣能換掉 0.99 的現任。

Q3 錯誤診斷

同事把訓練資產搬進新專案,執行結果一切正常,但 MLflow 上什麼都沒有。以下是他的資產與執行結果。最可能的原因?

@dg.asset def trained_model(context, mlflow_res: MlflowResource, train_test: dict) -> str: with mlflow.start_run(run_name="train") as run: # 直接就開 run ... # 執行結果 Dagster success = True # 沒有任何錯誤 管線的 tracking 裡有幾個 run: 0 工作目錄多出來的東西: ['mlflow.db']

資源被注入只代表「你拿到了設定」,套不套用是你的事——少了 setup()set_tracking_uri 就沒被呼叫,MLflow 用預設位置在當下的工作目錄開了一個新的 mlflow.db,訓練確實跑了、也真的有記錄,只是記到沒人會看的地方。這是最典型的「安靜地做錯」:Dagster 顯示成功、不會有任何紅字。A 不成立,資源沒註冊會在組圖時就報 DagsterInvalidDefinitionError: resource with key 'mlflow_res' ... was not provided,根本不會執行到;B 是杜撰的行為,MLflow 寫不進去會直接拋例外;C 症狀不符——掉進 Default 實驗的話,run 還是會出現在同一個 tracking 裡,而不是在別的資料庫。

Q4 情境題

每天凌晨的重訓連續三天被閘門擋下,AUC 從 0.97 掉到 0.86,而模型設定完全沒動過。線上還是上週那版 champion。你明天早上第一件該做的事是?

設定沒變、分數掉了,變的只可能是資料——這正是課堂上 run 4 的情境(同一組參數,只加了 drift,AUC 就從 0.9698 掉到 0.8641)。閘門已經幫你做完該做的事:擋下爛模型、保住線上那版、而且把「這次考幾分」留在紀錄裡;接下來要修的是資料,不是門檻。B 和 C 都是把警報關掉——問題還在,只是你看不到了,而且爛模型會直接上線。D 更糟:在不知道原因的情況下換模型,等於用更複雜的東西去硬記已經壞掉的資料,上線後會以更難察覺的方式失敗。

Q5 情境題

上游團隊每天不定時丟 1~5 批新資料進物件儲存,你希望「有新資料就重訓一次、同一批不要重跑」。最合適的觸發方式是?

「有事情發生才跑」正是感測器的定義,而 cursor 是它的記憶:記住上次看到哪一批,下一 tick 只處理新的;run_key 再幫你去重,同一批資料不會開出第二次執行。A 每 10 分鐘無條件重訓,一天 144 次訓練+144 次評估,浪費算力也把 MLflow 灌滿沒人看的 run,而且新資料最壞還是要等 10 分鐘。C 把自動化退回人工,你等於用一支 API 換來一個待辦事項。D 讓資產永遠不會結束,執行卡住、失敗也無法重試——輪詢是感測器的工作,不是資產的。

HANDS-ON · MOLAB

實作在 molab 跑(免費)

molab 的登入狀態進不了內嵌框架(瀏覽器的跨站 cookie 保護), 所以 notebook 要在新分頁執行——把它跟本頁並排開,左邊教學照樣對照。

  1. 登入 molab(GitHub / Google)
  2. 開啟課程 notebook,Fork 成自己的副本即可編輯
  3. 從第一格往下全部執行(首次安裝套件約 1–2 分鐘,管線本身合計約 1 分鐘上下)——免費 CPU 環境即可,不需要 GPU;MLflow 與 Dagster 的帳本都在暫存資料夾裡

不想用 molab?下載 mlops-pipeline_ext.py 後在自己電腦 uvx marimo edit --sandbox mlops-pipeline_ext.py,依賴會自動安裝。

molab 的線上編輯器在手機上體驗有限——動手這一段建議用電腦進行。