AI 互動教室 ‹ MLOps 自動化技術
下載 .py 開啟實戰 notebook ↗ 留言回報
DAGSTER · SOFTWARE-DEFINED ASSETS · 03

Dagster 軟體定義資產:
管線不是一串任務,是一張資料地圖

上一課你把最好的模型放進 Registry;但那份訓練資料是誰在什麼時候算出來的? 凌晨兩點的 cron 跑了三支腳本,模型今天預測怪怪的——你查不到特徵表是哪一批資料算的、清資料丟掉了幾列, 也沒辦法「只重算特徵、不重抓原始資料」。Dagster 把主角從任務換成資產: 每一份資料就是一個函式,依賴自己連成一張圖。先點點看——

點任何一個資產=對它按 Materialize

列數、客戶數與錯誤訊息都是 notebook 的實測結果(同一組亂數種子):清完 478 列、59 位客戶; 把 min_amount 拉到 800 只剩 106 列,低於品質閘門的 400 列底線。

01 · 資產

一個函式,就是一份資料

import dagster as dg @dg.asset(description="模擬的原始訂單(含退款負值)", group_name="raw") def raw_orders() -> pd.DataFrame: # 函式名 = 資產名 ... return pd.DataFrame({"order_id": ..., "customer": ..., "amount": ..., "returned": ...}) res = dg.materialize([raw_orders]) # 真的去算、存起來、記下一筆 materialization 事件 res.asset_value("raw_orders") # 500 列,其中 22 筆金額為負

資產(asset)是這一課唯一要記住的詞:一份「存在於某處、有人會用」的東西 (一張表、一個模型檔、一份報表、一個向量索引),加上「它是怎麼算出來的」那個函式。 腳本世界裡這兩件事是分開的——檔案在硬碟上,產生它的程式在某支 .py 的某幾行; Dagster 要求你把它們綁在一起宣告。

宣告不會產生資料,就像寫好食譜不等於做出菜。實體化(materialize)才是真的執行函式、 把結果交給儲存層、並留下一筆事件(時間、run id、中繼資料)。正式部署時你不會自己呼叫 materialize()——UI 上按 Materialize、排程、感測器都會做這件事; notebook 裡直接呼叫,是為了讓每一步看得見。

到 notebook 的 1️⃣ 節:第一個資產與 materialize
02 · 依賴成圖

參數名就是上游名:管線是推出來的,不是寫出來的

@dg.asset(group_name="clean") def clean_orders(context, config: CleanConfig, raw_orders: pd.DataFrame) -> pd.DataFrame: ... # 參數名 raw_orders = 上游資產名 @dg.asset(group_name="features") def customer_features(clean_orders: pd.DataFrame) -> pd.DataFrame: ... dg.materialize([raw_orders, clean_orders, customer_features]) # 執行順序:raw_orders → clean_orders → customer_features(沒有人寫過這行順序)

Airflow 那類工具要你手寫 a >> b >> c。管線小的時候沒問題; 等到四十張表互相引用,那串箭頭遲早跟真實依賴對不上——有人加了一張表忘了接線,某天它就悄悄用到昨天的資料。 Dagster 不讓你手寫順序:圖是從程式碼推出來的,程式怎麼寫、圖就長什麼樣,不會分岔。 實測這條線清完剩 478 列(丟掉 22 筆退款)、彙總出 59 位客戶。

既然名字就是契約,打錯一個字會怎樣?故意把 clean_orders 寫成 clean_order,Dagster 在組圖的時候就擋下來——還沒有執行任何函式、 沒有半筆資料被寫出去,而且它會猜你要的是哪個:

DagsterInvalidDefinitionError: Input asset "["clean_order"]" is not produced by any of the provided asset ops and is not one of the provided sources. Did you mean one of the following? ["clean_orders"]
到 notebook 的 2️⃣ 節:依賴成圖、血緣圖、名字打錯
03 · 中繼資料

每一次執行都留下隨身紀錄

context.log.info(f"kept {len(df)} / {len(raw_orders)} rows") # 日誌:給人當下讀,事後會被沖掉 context.add_output_metadata({ # 中繼資料:永久跟著這次實體化 "rows": len(df), "dropped": len(raw_orders) - len(df), "preview": dg.MetadataValue.md(df.head(3).to_markdown(index=False)), })

「清資料那步是不是把太多列丟掉了?」——要回答這種問題,光有「跑成功了」不夠。 中繼資料是跟著那一次實體化存下來的小抄:數字、字串、markdown、JSON、URL、檔案路徑都可以。 數字類在 UI 上會自動畫成時間序列,「今天的列數突然剩一半」一眼就看得到,不用等模型爛掉才回頭查。 實測這次記下 rows=478dropped=22 與前三列預覽。

到 notebook 的 3️⃣ 節:把 materialization 事件讀回來
04 · DEPS

只排順序、不傳資料——以及它最容易踩的坑

@dg.asset(deps=[customer_features], group_name="features") # 傳函式物件,不是字串 def feature_report() -> dg.MaterializeResult: # 真實世界這裡會去讀 warehouse、產 PDF、寄信 return dg.MaterializeResult(metadata={"recipients": 3, "status": "sent (simulated)"})

寫成參數的依賴同時做了兩件事:排順序把上游的值搬進來。但「把報表寄出去」只要等 customer_features 算完就好,它自己會去讀資料倉儲;上游是一張 300 GB 的表時, 你更不會想把它搬進 Python 記憶體。deps=[...] 就是這種依賴:有順序、有血緣、不傳值, 函式簽名裡不會出現上游的名字,通常回傳 MaterializeResult——「我做完了,這是我的紀錄」。

坑在這裡:deps 也吃字串,而字串打錯不會報錯。 在 Dagster 眼裡 "clean_order" 是一個合法的資產名,只是這份定義裡沒人負責算它—— 它變成一個外部資產(別的團隊、別的工具產生的)。實測結果:

@dg.asset(deps=["clean_order"]) # 少一個 s,Dagster 不會擋 def mail_report() -> dg.MaterializeResult: ... res = dg.materialize([raw_orders, clean_orders, mail_report]) res.success # True ← 什麼都沒壞 [ev.asset_key.to_user_string() for ev in res.get_asset_materialization_events()] # ['mail_report', 'raw_orders', 'clean_orders'] ← 報表比資料還早跑完

這是最難抓的那種 bug:run 是綠的、沒有任何訊息,只是報表用到的資料比你以為的舊。 能傳函式物件就別傳字串——打錯字時 Python 自己就會 NameError,連 Dagster 都不用出手。 真的要用字串(跨檔案、上游是別人的資產),就在 Definitions 建好後檢查圖上有沒有意料之外的外部資產。

到 notebook 的 4️⃣ 節:deps 與打錯字的血緣圖
05 · IO MANAGER

「算什麼」與「存哪裡」分家

class CsvIOManager(dg.ConfigurableIOManager): root: str def handle_output(self, context, obj): # 資產算完 → 存 obj.to_csv(self._path(context), index=False) def load_input(self, context): # 下游要用 → 讀 return pd.read_csv(self._path(context)) dg.materialize([raw_orders, clean_orders, customer_features], resources={"io_manager": CsvIOManager(root=CSV_ROOT)}) # 只換這一行

你的資產函式裡沒有一行 to_csvread_csv, 那 clean_orders 的 DataFrame 是誰存的、又是誰讀給下游的? 答案是 IO manager。預設的 fs_io_manager 把每個資產 pickle 成一個檔案 (檔名=資產名,沒有副檔名);換成上面這個 CSV 版之後,落地變成 ['clean_orders.csv', 'customer_features.csv', 'raw_orders.csv']—— 三個資產函式一個字都沒改。今天存本機、明天上 S3、後天寫進 Snowflake,都是換這一個元件的事。

selection=[customer_features]只算這一個;上游 clean_orders 由 IO manager 從上次的結果載回來(事件裡看得到 1 筆 LOADED_INPUT)。
selection="clean_orders*"它自己+所有下游,實測實體化 ['clean_orders', 'customer_features']raw_orders 不用重抓。

這就是開頭那個「可以嗎」的答案:可以,前提是上游真的在這個 storage 裡算過。 換一台機器、換一個 storage 目錄、或那個資產從來沒被實體化過,載入就會失敗—— 這是新手最常撞的一堵牆(本課測驗會再考一次)。

到 notebook 的 5️⃣ 節:換 IO manager、只重算下游
06 · ASSET CHECK

資料不對,就不要拿去訓模型

@dg.asset_check(asset=clean_orders, description="清完不該再有負金額") # 不 blocking def no_negative_amount(clean_orders: pd.DataFrame) -> dg.AssetCheckResult: bad = int((clean_orders["amount"] < 0).sum()) return dg.AssetCheckResult(passed=bad == 0, severity=dg.AssetCheckSeverity.WARN, metadata={"bad_rows": bad}) @dg.asset_check(asset=clean_orders, blocking=True) # 這個會擋下游 def enough_rows(clean_orders: pd.DataFrame) -> dg.AssetCheckResult: n = len(clean_orders) return dg.AssetCheckResult(passed=n >= 400, severity=dg.AssetCheckSeverity.ERROR, metadata={"rows": n, "min_rows": 400}) dg.materialize([raw_orders, clean_orders, customer_features, no_negative_amount, enough_rows]) # 檢查跟資產放同一個清單

上游今天改了一個欄位定義,你的清理邏輯照跑不誤,只是留下的列數掉了一大半。程式沒拋例外、run 是綠的、 模型照樣訓完上線——三天後才有人發現預測全歪。資產檢查就是把「資料應該長什麼樣」寫成程式碼掛在資產上: 像單元測試,但測的是資料,而且每次實體化後自動跑,結果跟資產一起記錄。

兩個常被混在一起的旋鈕其實是分開的:severity(有多嚴重,寫在結果上,預設就是 ERROR, 想降級寫 AssetCheckSeverity.WARN)只影響這筆紀錄長什麼樣; blocking=True(寫在裝飾器上)才是真的踩煞車。 實測把 min_amount 拉到 800,只剩 106 列(門檻拉到 250 就守不住 400 列底線):

res.success # False [ev.asset_key.to_user_string() for ev in res.get_asset_materialization_events()] # ['raw_orders', 'clean_orders'] ← 沒有 customer_features,下游真的被擋住了 dagster._core.errors.DagsterAssetCheckFailedError: 1 blocking asset check failed with ERROR severity: clean_orders: enough_rows

no_negative_amount 這次照樣執行、照樣通過——它不 blocking,就算失敗也只是留下一筆紀錄。 資產函式自己丟例外(上游 schema 變了、API 掛了)則是另一種失敗:Dagster 把你的例外包起來, 外層是 DagsterExecutionStepExecutionError,原始的 ValueError: boom: upstream schema changed 收在 error.cause 裡; 下游一樣不跑,而且該資產在圖上維持「上一次成功的樣子」,不會被半成品覆蓋。

到 notebook 的 6️⃣ 節:閘門關上的那一刻+拉桿互動
07 · DEFINITIONS

收成一份,交給 dagster dev

defs = dg.Definitions( assets=[raw_orders, clean_orders, customer_features, feature_report], asset_checks=[no_negative_amount, enough_rows], resources={"io_manager": CsvIOManager(root=CSV_ROOT)}, ) # $ dagster dev -f my_pipeline.py → http://localhost:3000

Definitions 是「這個專案有哪些東西」的唯一入口:資產、檢查、資源, 還有下一課的排程與感測器。dagster dev 起來之後,UI 會畫出你在 notebook 裡看到的同一張圖, 另外標上每個資產的最近實體化時間、中繼資料趨勢、檢查狀態;你可以點任何一個資產按 Materialize、只重算某個子集、 翻每次 run 的日誌。

它也是整體檢查的地方:資產名全域唯一,撞名不是警告而是錯誤——否則「這份資料是誰算的」會有兩個答案。 實測訊息:DagsterInvalidDefinitionError: Duplicate asset key: AssetKey(['clean_orders'])。 真實專案常用 AssetKey(["marketing", "clean_orders"]) 這種多層命名把不同來源分開。

到 notebook 的 7️⃣ 節:Definitions 與撞名
08 · 實戰

換你動手

LEVEL 1

加一個資產 top_customers:吃 customer_features,回傳 total 最高的 5 位。實體化整條線,確認它出現在血緣圖上、中繼資料記了 5 列。

LEVEL 2

customer_features 加一個 blocking=True 的檢查:return_rate 必須在 0–1 之間且沒有 NaN。故意弄壞一筆資料,看它擋不擋得住下游。

LEVEL 3

CsvIOManager 改成依型別決定格式:DataFrame 存 CSV、其他型別存 pickle。驗證:回傳 MaterializeResultfeature_report 根本不會經過 IO manager。

卡住了?每一題在 notebook 末節都有折疊解答——先自己做,再打開對照。

09 · 驗收

情境測驗

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

Q1 情境題

目前三支腳本靠 cron 串起來:2:00 抓訂單、2:10 清資料、2:20 算特徵。今天特徵表怪怪的,你想知道它是哪一批資料算的、也想只重算特徵。最該做的第一步是?

問題的根源是「管線記錄的是跑了哪些腳本,不是產出了哪些資料」。宣告成資產之後,三件事一次到位:依賴自動成圖(不用手寫順序)、每次實體化留下時間與中繼資料(列數就在那裡)、可以只選一個資產重算而上游從儲存層載回。A 加了日誌但沒有結構,日誌會被沖掉、也還是不能只重算一段;B 只解決順序,「那張表是什麼時候、用什麼算的」依然無解;D 是手工版的血緣,量一大就沒人維護,而且沒有回答「誰算的」。

Q2 錯誤診斷

報表資產每天都「成功」,但業務說報表數字比實際慢一天。程式與執行結果如下,最可能的原因是?

@dg.asset(deps=["clean_order"]) # 報表:等資料清好再寄 def mail_report() -> dg.MaterializeResult: ... >>> res = dg.materialize([raw_orders, clean_orders, mail_report]) >>> res.success True >>> [ev.asset_key.to_user_string() for ev in res.get_asset_materialization_events()] ['mail_report', 'raw_orders', 'clean_orders']

關鍵線索在執行順序:報表跑在 clean_orders 之前(這次甚至是第一個),而 run 完全沒有報錯。deps 接受字串,而任何字串在 Dagster 眼裡都是一個合法的資產名——打錯的 clean_order 變成一個外部資產(圖上有節點、沒人負責算),報表於是「沒有任何上游要等」。修法是傳函式物件 deps=[clean_orders],打錯字時 Python 直接 NameError。A 說反了:deps 只要名字對就會正確排序,這裡是名字不對;C 症狀相似但原因不同,這次 clean_orders 確實有重算,只是排在報表後面;D 是常見誤解——清單順序不影響執行順序,順序完全由依賴推導。

Q3 情境題

上游偶爾會送來只有零星幾列的殘缺資料。你不希望模型拿這種資料去訓練,但也不想因此讓整條管線變得難維護。最佳做法是?

blocking 的資產檢查正是為這件事設計的:條件寫在資產旁邊(誰在管這份資料的品質一目了然)、每次實體化自動跑、不合格時下游根本不會開始,而且結果會留在 UI 上有歷史可查。B 能擋住但把資料品質的判斷藏進訓練程式裡,換一個下游就要再寫一次,UI 上也只會看到「訓練失敗」而不是「資料不合格」;C 只是提醒,信件寄出去的時候模型早就用爛資料訓完了;D 最危險——它讓管線「看起來一直是好的」,血緣紀錄還會顯示今天算過,實際上是舊資料,屬於那種三個月後才爆的錯。

Q4 錯誤診斷

同事在自己電腦第一次跑這個專案,想直接算最下游的特徵表,結果如下。最可能的原因與修法?

res = dg.materialize([raw_orders, clean_orders, customer_features], selection=[customer_features], instance=my_instance) dagster._core.errors.DagsterExecutionLoadInputError: Error occurred while loading input "clean_orders" of step "customer_features": FileNotFoundError: [Errno 2] No such file or directory: '/tmp/tmp7ol1w4mr/storage/clean_orders'

「只重算下游」的前提是上游的結果已經存在這個 instance 的 storage 裡;全新的環境什麼都沒有,IO manager 去讀 storage/clean_orders 自然找不到檔案。訊息的兩層說得很清楚:外層是「載入 customer_features 的輸入 clean_orders 時出錯」,底層是 FileNotFoundError。Dagster 不會自動幫你補跑上游——它只做你叫它做的事。A 的症狀不同:參數名對不上會在組圖時就被擋下,訊息是 Input asset ... is not produced by any of the provided asset ops,而且根本跑不到載入這一步;B 誤讀訊息,是「檔案不存在」不是「不能寫」;C 也錯,檢查失敗的訊息是 DagsterAssetCheckFailedError,長得完全不一樣。

Q5 情境題

資料工程團隊要求:所有中間資料從本機 pickle 改存成公司資料湖上的 Parquet。你的專案有 30 個資產。工作量最小、風險最低的做法是?

「算什麼」與「存哪裡」分家就是為了這一天:IO manager 只要回答 handle_outputload_input 兩個問題,換掉它,30 個資產函式一個字都不用改,也不會有人漏改。A 把儲存細節塞回每個函式裡,下游還要自己知道路徑,等於放棄了 Dagster 幫你管的那一層,下次再換格式又要改 30 個地方;C 讓資產數量翻倍、血緣圖被搬運節點灌爆,而且真正的資料還是躺在 pickle 裡;D 是事後補救,管線執行當下用的仍是舊格式,轉檔腳本一失敗就出現「兩份不同步的資料」。

HANDS-ON · MOLAB

實作在 molab 跑(免費)

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

  1. 登入 molab(GitHub / Google)
  2. 開啟課程 notebook,Fork 成自己的副本即可編輯
  3. 從第一格往下全部執行(首次安裝套件約 1–2 分鐘)——免費 CPU 環境即可,不需要 GPU;全部在暫存資料夾裡跑,不連任何伺服器

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

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