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

Dagster 自動化:
誰來按下那個「執行」?

上一課你宣告了資產、看著 Dagster 從依賴推出一張圖——然後自己呼叫了 materialize()。 真實世界不會有人每天凌晨 2 點坐在電腦前按按鈕。這一課補上被留下的那半:誰按、什麼時候按、按下去跑哪一段、壞了誰重試。 先玩最基本的那個問題——「這個 cron 到底什麼時候會跑」,順便撞一次幾乎每個人都踩過的時區陷阱:

接下來 5 次觸發是在瀏覽器裡照 cron 規則算的;時區欄位的行為與錯誤訊息來自 notebook 實測—— 排程沒寫 execution_timezone 時那一欄真的是 None(以 UTC 解讀), cron 寫錯時 Dagster 在建立排程的當下就丟出那句話。

01 · 環境、參數、打包

資產要被自動觸發,先把三件事分清楚

自動化不是「把 materialize() 塞進 cron」。要讓別人(排程、感測器、daemon)替你按執行, 程式得先能被從外面設定:連線與路徑歸資源,這一跑的參數歸設定,要一起跑的資產打包成 job。

class FeatureStore(dg.ConfigurableResource): # resource:外部世界的接點,部署時決定(dev / prod 各一份) root: str class IngestConfig(dg.Config): # config:這一次執行的參數,每次觸發都可以不同 n_rows: int = 500 seed: int = 0 @dg.asset def orders(context, config: IngestConfig, store: FeatureStore) -> pd.DataFrame: ... # 寫成參數就會被餵進來,跟上游資產一樣 nightly_job = dg.define_asset_job("nightly_job", selection=dg.AssetSelection.assets("orders").downstream())

實測:同一個 orders 函式跑兩次,程式碼一個字沒改——第一次配 dev 資源抓 20 筆、 第二次配 prod 資源抓 500 筆,檔案分別落在兩個目錄。之後排程與感測器要「帶著參數觸發」,帶的就是 {"ops": {"orders": {"config": {…}}}} 這個字典:它們不是呼叫你的函式,是遞一張寫好參數的單子

job 用 AssetSelection 選資產的好處是選擇會自己長大:實測跑 nightly_job 實體化了 ['orders', 'orders_report']——我們沒點過 orders_report 的名字, 是 .downstream() 把它帶進來的。之後新增的下游資產也一樣,半夜那一跑自動包含它。

到 notebook 的 1️⃣–2️⃣ 節:資源與設定、job 真的跑一次
02 · 排程

時間到就跑:cron、時區,以及那張「單子」

nightly_schedule = dg.ScheduleDefinition( name="nightly_2am", job=nightly_job, cron_schedule="0 2 * * *", # 分 時 日 月 週 execution_timezone="Asia/Taipei", # ← 不寫的話這一欄是 None,Dagster 以 UTC 解讀 ) @dg.schedule(job=nightly_job, cron_schedule="0 2 * * *", execution_timezone="Asia/Taipei") def nightly_sized(context): # 想「依日期決定參數」就自己寫函式 day = context.scheduled_execution_time.strftime("%Y-%m-%d") return dg.RunRequest( run_key=day, # 同一天只會開一次 run run_config={"ops": {"orders": {"config": {"n_rows": 200}}}}, # 帶參數 tags={"day": day, "trigger": "schedule"}, # 貼標籤 )

排程的產出不是「執行」,是一張 RunRequest(請跑這個 job)。 在 notebook 裡不用起 daemon 也能看它:evaluate_tick(build_schedule_context(scheduled_execution_time=…)) 直接問「這個時刻你會發什麼」。實測最陽春的排程送出 1 張,內容是 run_key=Nonerun_config={}tags={'dagster/schedule_name': 'nightly_2am'}(Dagster 自動貼的)。

run_key 是防重複跑的關鍵:同一個排程送出相同 run_key 的單子,Dagster 只開一次 run—— daemon 重評估、服務重啟、時鐘回撥都不會害你算兩次帳。日期字串是最常見的 run_key。

cron 寫錯不會安靜地失敗,建立排程物件的當下就爆(所以 dagster dev 一開就會告訴你)。實測原文:

DagsterInvalidDefinitionError: Found invalid cron schedule '0 25 * * *' for schedule 'typo''. Dagster recognizes standard cron expressions consisting of 5 fields.
到 notebook 的 3️⃣ 節:兩種排程寫法、tick 送出的單子、cron 驗證
03 · 分割與補跑

把資產切成一天一片,缺哪幾天看得見

「一整包」的資產只有兩種狀態:跑過、沒跑過。真實資料是一天一天長出來的,所以你會需要 只重算昨天那一片補跑掛掉那三天分天比較品質——這些在一整包上都做不到。

daily_parts = dg.DailyPartitionsDefinition(start_date="2026-09-01") @dg.asset(partitions_def=daily_parts) def daily_orders(context) -> pd.DataFrame: day = context.partition_key # 這一跑負責哪一片 ... dg.materialize([daily_orders], partition_key="2026-09-03", instance=inst) instance.get_materialized_partitions(dg.AssetKey("daily_orders")) # 哪幾片有了 → 剩下的就是要補的 daily_job = dg.define_asset_job("daily_orders_job", selection=[daily_orders], partitions_def=daily_parts) daily_schedule = dg.build_schedule_from_partitioned_job(daily_job, hour_of_day=3) # 每天 03:00 跑「前一天」那片

notebook 用「今天往前 7 天」當分割起點,所以你哪天執行都會有 7 片。故意讓排程只跑了最早 3 片 (模擬掛掉),get_materialized_partitions 立刻列出缺的 4 片,一個迴圈補完變 7 / 7—— 每一片都是獨立的一次 run、各自留下 rowstotal 中繼資料,直接拿來畫成每日趨勢。

build_schedule_from_partitioned_job 實測產出 cron 0 3 * * *、 時區 UTC(又是那個預設,記得改),而且它送出的單子帶的是 partition_key:任何一天的 03:00 tick 都指向前一天那一片。 「每天凌晨結算昨天」不用自己算日期。

分割也有專屬的錯誤:忘了給 partition_keyDagsterInvariantViolationError: Cannot access partition_key for a non-partitioned run, 給了範圍外的日期是 DagsterUnknownPartitionError: Could not find a partition with key ...

到 notebook 的 4️⃣ 節:補跑、每日趨勢圖、挑日期問排程
04 · 感測器

有事情發生才跑,靠 cursor 記住看過什麼

排程回答「時間到了嗎」,但客戶什麼時候上傳檔案、上游什麼時候補完資料,都不看時鐘。 感測器就是那個「每 30 秒被叫起來問一次」的函式,它只能回答兩件事: 要跑什麼(RunRequest,或為什麼不跑(SkipReason

@dg.sensor(job=nightly_job, minimum_interval_seconds=30) def inbox_sensor(context): seen = set(json.loads(context.cursor)) if context.cursor else set() # cursor 是這個函式唯一的記憶 files = sorted(p.name for p in INBOX.glob("*.csv")) new = [f for f in files if f not in seen] for f in new: yield dg.RunRequest(run_key=f, tags={"file": f}) # 同一個檔案只會開一次 run if not new: yield dg.SkipReason(f"no new files (已看過 {len(seen)} 個)") context.update_cursor(json.dumps(sorted(seen | set(new)))) # 把記憶存回去

實測四次 tick:空資料夾 → skip;丟進兩個檔案 → 2 張單子,cursor 變成 ["batch_a.csv", "batch_b.csv"];再問一次 → skip(檔案還在,但看過了); 再來一個新檔 → 1 張。cursor 存什麼由你決定:檔名清單、上次的最大 mtime、上次讀到的資料庫 id。

盯自己人則用 資產感測器@dg.asset_sensor(asset_key=…, job=…) 幫你把 cursor 那段寫好了—— 上游資產一有新的 materialization 就觸發。實測它的 cursor 從 None 變成一個數字 (事件的 storage id,也就是「我讀到第幾筆事件」),沒有新事件時的 skip 訊息是 No new materialization events found for asset key AssetKey(['orders'])

到 notebook 的 5️⃣ 節:三種 tick、資產感測器、自己排一次 tick
05 · 宣告式自動化

不寫 job,資產自己說什麼時候該更新

排程與感測器都是由外往內推:這個時刻/這件事發生時,去跑那包資產。管線一大就變成 20 個排程、 8 個感測器,每加一個資產都要想「它該掛在哪個 job 上」。宣告式自動化反過來:條件寫在資產上。

@dg.asset(automation_condition=dg.AutomationCondition.eager()) # 上游一更新,我就跟著更新 def orders_alert(): ... @dg.asset(automation_condition=dg.AutomationCondition.on_cron("0 6 * * *")) # 每天 6 點後、等上游備妥才更新 def daily_report(): ... result = dg.evaluate_automation_conditions(defs=defs, instance=inst, cursor=前一次的cursor) result.total_requested, result.get_num_requested(dg.AssetKey("orders_alert")) # 沒有 run_requests 屬性

實測 eager() 的三次評估(每次把上一次的 cursor 傳進去): 什麼都還沒實體化 → 0;上游 orders 剛跑完 → 1(下游被請求更新); 下游也跑完後再評 → 0全程沒有寫任何 job、任何排程,正式環境那一個「1」會由 daemon 直接變成一個 run。

最容易誤會的是 on_cron:它不是「每天 6 點跑我」,而是 「每天 6 點之後,等所有上游在這個週期內更新過了,才跑我」。notebook 把時鐘捏在手上跑了六步, 只有第五步被請求:cron 時刻早就過了,但上游是在那之後才更新的。 差別正是資料工程最常見的競態——排程時間到了、上游還沒好,於是你用舊資料算出一份新報表。 排程不會等上游,on_cron 會。

schedule時鐘觸發,要 job。適合「每天固定時間」的入口資產。陷阱:時區預設 UTC、不等上游。
sensor你寫的檢查邏輯,要 job。適合「外面發生了什麼」。陷阱:忘了更新 cursor、忘了給 run_key。
AutomationCondition資產自己的條件,不用 job。適合中下游那一大片。陷阱:以為 on_cron 是排程。

實務上三種混用:入口資產(要去外面撈資料的那幾個)用排程或感測器,中下游用 eager() 自己跟上——要維護的排程只有幾個,不是幾十個。

到 notebook 的 6️⃣ 節:eager 三次評估、on_cron 六步實驗
06 · 失敗與收成

壞掉會重試、有人被通知,然後交給 daemon

自動化最現實的一件事:沒有人在看的時候,東西一定會壞。網路抖一下、連線被回收、上游 API 回 503—— 這種「再試一次就好」的失敗,不該讓整條管線停到早上。

@dg.asset(retry_policy=dg.RetryPolicy(max_retries=2, delay=0.2)) def flaky_train(context) -> int: if context.retry_number < 2: # 第 0、1 次故意失敗 raise RuntimeError("connection reset") return 42 @dg.run_failure_sensor # 放進 Definitions 的 sensors,由 daemon 監看 def alert_on_failure(context): send_slack(context.failure_event.message)

實測事件序列:STEP_START → STEP_UP_FOR_RETRY → STEP_RESTARTED → STEP_UP_FOR_RETRY → STEP_RESTARTED → STEP_SUCCESS, run 成功、資產的值是 42——綠燈,但你看得到它抖過。所以不要用 try/except 把暫時性失敗吞掉,那會讓你以為系統很健康。 重試用完還是失敗時,Dagster 的最後一句話是 Exceeded max_retries of 1,這時該被叫起來的是 run failure sensor。 (資料不對是另一回事:那是上一課的資產檢查,不該重試。)

defs = dg.Definitions(assets=[...], jobs=[...], schedules=[...], sensors=[...], resources={...}) $ dagster dev -f pipeline.py # 開 http://localhost:3000,同時起 webserver 與 daemon

最後把零件收成一份 Definitions(實測這一課的成品:4 個資產、3 個 job、3 個排程、3 個感測器, Dagster 另外自動補一個涵蓋全部資產的隱含 job __ASSET_JOB)。 daemon 是這一課的隱形主角:它一直醒著,看 cron 到了沒、每 30 秒問一次每個感測器、評估所有資產的自動化條件、 監看 run 狀態觸發失敗通知。UI 上每個排程與感測器都有一個開關(預設是關的),也看得到每一次 tick 發了幾張單、skip 的理由是什麼—— 本課用 evaluate_tick 看到的東西,就是那些 tick 紀錄的內容。

到 notebook 的 7️⃣–8️⃣ 節:重試事件、失敗通知、Definitions 全景圖
07 · 實戰

換你動手

LEVEL 1

寫一個「每週一早上 9 點」的排程,帶 run_confign_rows 設成 1000,用 evaluate_tick 確認送出的 run_key 與參數是你要的。

LEVEL 2

把檔案感測器的 cursor 從「檔名清單」改成「上次看到的最大 mtime」,並限制一次最多送 2 張單子。丟 5 個檔案、連跑三次 tick 驗證行為。

LEVEL 3

把分割資產接上 AutomationCondition.on_cron("0 3 * * *"),再加一個吃它的下游用 eager(),串著 cursor 評估,觀察哪幾片被請求(get_requested_partitions)。

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

08 · 驗收

情境測驗

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

Q1 情境題

客戶會不定時把 CSV 丟進一個共用目錄,一天可能 0 次也可能 20 次,你要在檔案出現後幾分鐘內處理它,而且同一個檔案只能處理一次。最佳做法是?

感測器就是為這種「不看時鐘、看事件」的觸發而生的:cursor 記住看過什麼(不用每次重掃全部),run_key 是最後一道保險——同一個 run_key 的單子 Dagster 只會開一次 run。A 每 5 分鐘把整個目錄重跑一次,處理過的檔案會一再被處理,量大時成本很可觀;B 的全域變數活在 daemon 的記憶體裡,daemon 一重啟或換一台機器就全忘了,cursor 存在 Dagster 的儲存體裡才不會;D 誤會了 eager()——它看的是「上游資產有沒有新的實體化」,不會去看目錄裡有沒有新檔案,外部世界的事件還是要靠感測器帶進來。

Q2 錯誤診斷

你要一個「每天凌晨 2 點」的排程,寫好上線後發現它每天都在早上 10 點才跑。你去 REPL 確認了排程物件:

>>> sched = dg.ScheduleDefinition(job=nightly_job, cron_schedule="0 2 * * *") >>> sched.execution_timezone None

差整整 8 小時就是時區的簽名——台北是 UTC+8。execution_timezone 沒設定時實測是 None,Dagster 就以 UTC 解讀那串 cron;同樣的事情也發生在 build_schedule_from_partitioned_job 產生的排程上(它的時區印出來是 UTC)。A 的欄位順序是「分 時 日 月 週」,"0 2 * * *" 沒寫錯;C 不成立,daemon 是持續執行的行程,它啟動時會補評估,不會把每天的時刻整體推遲固定 8 小時;D 是併發限制的症狀(run 排隊),不會讓觸發時間每天穩定晚 8 小時。

Q3 錯誤診斷

你想要「每小時的第 25 分鐘跑一次」,寫成 cron_schedule="0 25 * * *",結果 dagster dev 一開就失敗:

DagsterInvalidDefinitionError: Found invalid cron schedule '0 25 * * *' for schedule 'typo''. Dagster recognizes standard cron expressions consisting of 5 fields.

把 cron 看成「時 分」是最常見的誤記——第一欄是、第二欄才是,所以 0 25 * * * 等於「每天 25 點整」,25 超出 0–23 就被擋下。要「每小時第 25 分」是 25 * * * *。A 方向錯了,訊息說的正是「標準 5 欄位」;B 誤讀了訊息,'typo' 只是那個排程的名字,跟合不合法無關;C 剛好相反——這類錯誤是在建立排程物件的當下就丟出來的,dagster dev 載入定義就會失敗,它根本上不了線。

Q4 情境題

你的報表每天早上 6 點要出,但它依賴的訂單資料由另一個團隊的管線在早上 5 點到 7 點之間跑完,時間不固定。你希望報表一定用當天的新資料,而且一天只出一次。最佳做法是?

on_cron 正是為這種情況設計的:它不是「時間到就跑」,而是「這個 cron 週期內、等所有上游都更新過了才跑,而且一個週期最多跑一次」——notebook 的六步實驗看得很清楚,cron 時刻過了但上游還是上一個週期的資料時它不動,上游更新的那一刻才被請求。B 就是最經典的競態:上游 6 點 20 分才好,你 6 點整用昨天的資料出了一份新報表,而且看起來完全正常。C 只是把猜測的等待時間拉長,上游哪天慢了照樣中獎,還每天白等兩小時。D 保證用新資料,但上游一天更新三次它就跑三次,違反「一天一次」,而且完全不受 6 點這個業務時間約束。

Q5 情境題

每天凌晨的訂單管線因為上游 API 偶爾回 503 而失敗,重跑一次通常就好了。你不想每天早上手動重跑,也不想失敗被靜靜吞掉。最佳做法是?

RetryPolicy 就是為「暫時性故障」設計的:實測事件流會留下 STEP_UP_FOR_RETRY → STEP_RESTARTED,重試成功後 run 是綠的、但你看得到它抖過幾次;真的救不回來時(Exceeded max_retries of ...)run failure sensor 會把人叫起來。A 是最危險的選項——用舊資料冒充新資料,管線永遠綠燈,錯誤要等到有人發現報表不動了才爆;B 讓同一天的資料被算很多次(沒有 run_key 保護時尤其亂),而且沒解決「失敗沒人知道」;D 方向不同:資產檢查管的是「資料對不對」,那種失敗不該重試,而這裡的問題是「連線抖了一下」。

HANDS-ON · MOLAB

實作在 molab 跑(免費)

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

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

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

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