使用 MongoDB 變更串流熱抽換 Cron 排程
寫死的 cron 意味著每次排程變更時都必須重新部署。將排程儲存在資料庫中,並讓一個守護行程透過變更串流來回應編輯動作。以下是其設計與故障處理方式。
如果您的背景作業是根據寫死在程式碼或環境變數中的 cron 表達式來執行,那麼要變更作業的執行時間就意味著需要重新部署。將一份報告的執行時間從早上 9 點改到 8 點,您就得交付一個建置版本、執行一個管線,並重新啟動一個程序,而這一切只為了變更一個數字。這樣做很慢,其風險與變更不成比例,而且會阻礙那些能維持系統健康的微調。解決方法是將排程儲存在資料庫中,並讓排程器即時對編輯做出反應。MongoDB 變更串流能讓反應立即發生。這篇文章涵蓋了其設計、使其可靠的重新連線處理機制,以及何時這種機制會過於複雜,超出您的需求。
寫死在程式碼中的排程所帶來的問題
寫死的排程將兩件變化速率和變化原因完全不同的事情耦合在一起。工作的邏輯會隨著工作內容的改變而改變,這是一個屬於建置過程的工程決策。工作的執行時間會隨著營運需求而改變,例如報告應在會議前產出、繁重任務應移出尖峰時段,而這是一個與程式碼無關的營運決策。將它們綁在一起意味著每一次對時間的微調,都得背負著程式碼部署的沉重代價。
成本不僅僅是部署時間,還有寒蟬效應。當變更排程需要一次生產環境的部署時,人們就會停止這麼做。排程會逐漸與實際需求脫節,因為沒有人願意為了一個小時的任務調整而發布一個建置版本。排程應該是可以編輯的資料,而不是必須發布的程式碼。
將排程視為資料
第一步是將排程放入一個集合中,每個工作一個文件,並將 cron 運算式作為一個欄位。
// One document per scheduled job. The cron expression is data, not code.
type Schedule struct {
ID string `bson:"_id"` // e.g. "daily_report"
CronExpr string `bson:"cron_expr"` // e.g. "0 9 * * *"
UpdatedAt time.Time `bson:"updated_at"`
}
現在,變更排程只需一次資料庫更新,無需部署:
db.schedules.updateOne(
{ _id: "daily_report" },
{ $set: { cron_expr: "0 8 * * *", updated_at: new Date() } }
)
但僅有資料是不夠的。必須有某個東西注意到變更並重新設定正在運行的排程器。你有兩種方法可以實現這一點,而它們之間的差異正是這篇文章的重點。
輪詢與反應
最簡單的方法是輪詢:排程器每分鐘重新讀取集合,並在有任何變更時重建其計時器。這種方法可行,而且對許多系統來說已經完全足夠。它的弱點是延遲和浪費。一項變更最多需要一個完整的輪詢間隔才會生效,而絕大多數的輪詢都發現沒有任何變更,純粹是額外開銷。縮短間隔以減少延遲,你會增加浪費。這就像一個沒有好設定的旋鈕,只有比較不差的設定。
變更串流移除了這個旋鈕。排程器不再是透過計時器詢問「有什麼變動嗎?」,而是訂閱一個即時的資訊流,並在有變動的瞬間被告知。排程的編輯幾乎是即時傳播的,而當沒有任何變更時,排程器什麼也不做,也不會產生任何成本。你同時獲得了更低的延遲和更少的工作量,這點很不尋常,值得為這個機制付出額外的關注。
何謂變更串流
變更串流是一個即時、有序的饋送,記錄著發生在集合上的寫入操作。在底層,它讀取資料庫用來保持複本集成員同步的相同複寫日誌,因此它能依序看到插入、更新和刪除操作被提交,而無需輪詢。您的常駐程式只需開啟串流一次,然後就會阻塞,每當排程文件發生變更時,就會收到一個事件。
// Watch the schedules collection. Block until an edit arrives, then apply it.
stream, err := coll.Watch(ctx, mongo.Pipeline{})
if err != nil {
return fmt.Errorf("open change stream: %w", err)
}
defer stream.Close(ctx)
for stream.Next(ctx) {
var event struct {
FullDocument Schedule `bson:"fullDocument"`
}
if err := stream.Decode(&event); err != nil {
logger.Warn("decode change event", logger.Err(err))
continue
}
scheduler.Reload(event.FullDocument) // rebuild this job's timer
}
當操作員執行該單行更新時,常駐程式會在短時間內收到事件,並僅重建受影響的計時器。無需重新部署,無需重新啟動,也無需等待輪詢間隔。現在,排程是真正的即時資料。
變更串流需要一個複本集,因為它們讀取複寫日誌,而獨立伺服器並不會維護此日誌。任何生產環境的 MongoDB 部署都已經是複本集,所以這很少會是一個新的要求,但在您想於本地的單一節點設定上使用此功能前,值得了解它在該環境下是不可用的。
使其可靠的部分:恢復權杖
一個單純的變更串流會有間隙。如果連線因網路瞬斷、容錯移轉或短暫重啟而中斷,而你只是簡單地重新開啟串流,你會從「現在」這個時間點恢復,並默默地錯過所有在斷線期間發生的變更。在三十秒的重新連線期間編輯的排程將會遺失,而守護行程會繼續執行舊的時間設定,儘管資料庫中的內容已然不同。這種分歧正是整個設計旨在防止的錯誤。
變更串流用「恢復權杖」(resume token)來解決這個問題。每個事件都帶有一個權杖,標記其在資訊流中的位置。將你已處理的最新權杖持久化保存,並在斷線後,從該權杖而不是從「現在」開始重新開啟串流。資料庫會依序重播你錯過的每一個變更,讓你從上次中斷的地方準確地接續上。
// Resume from the last processed position so a reconnect misses nothing.
opts := options.ChangeStream()
if resumeToken != nil {
opts.SetResumeAfter(resumeToken)
}
stream, err := coll.Watch(ctx, mongo.Pipeline{}, opts)
// ... after handling each event, persist stream.ResumeToken() durably
要遵守的準則是,在處理完每個事件後儲存權杖,其儲存方式必須足夠持久,以應對守護行程的重啟。重新連線時,載入它並從中恢復。有了這個機制,串流不僅快速,而且在任何長時間執行的消費者中都不可避免的斷線情況下,也能做到無間隙。權杖是將一個便利的功能轉變為可靠功能的關鍵。
一個注意事項:恢復權杖僅在變更仍然存在於複寫日誌(replication log)中時才有效,而該日誌的大小是有限的。如果守護行程停機時間過長,導致最舊的未處理變更已經從日誌中被輪替掉,那麼權杖就會失效,串流也無法從中恢復。處理這種情況的方法是,偵測到無效恢復錯誤(invalid-resume error),並退回執行集合的完整重新載入,這會將排程器重新同步到當前的真實狀態。無論如何,這與你首次啟動時所期望的從頭開始核對的路徑是相同的。
套用變更,毫不間斷
接收事件只是工作的一半。套用變更時必須小心謹慎,因為您正在變動一個執行中的排程器。應針對已變更的工作專門重建計時器,而不是拆除並重新建立每個計時器,如此一來,對單一排程的編輯便不會干擾其他排程,或在交換期間冒著重複觸發的風險。取消該工作 ID 的舊計時器,從更新後的 cron 運算式安裝新的計時器,並讓所有其他工作保持不變。保持重新設定的冪等性,如此一來,即使套用同一個事件兩次(這種情況可能由 resume-from-token 在重新連線的邊界處引起),排程器最終的狀態也會與只套用一次時相同。這種冪等性正是確保已恢復串流的「至少一次」特性安全無虞的關鍵。
何時應保持簡單
當排程變更足夠頻繁,以至於重新部署成為真正的拖累,或者當變更生效的延遲確實很重要時,變更串流便是正確的工具。如果你的排程幾乎從不變更,那麼操作上的效益就很小,而一個簡單的重啟以重新載入,或是一個緩慢的輪詢,建構起來更簡單,也更不容易出錯。不要為了一個一年只會有人編輯兩次排程的系統,加入即時事件消費者、恢復權杖持久化,以及無效權杖的備援機制。
更深層的原則可推廣到 cron 之外。任何因操作原因而非程式碼原因變更的設定,都應該存放在執行中系統所監看的資料儲存區中,這樣一來,變更它就只是一次編輯,而不是一次發佈。排程是一個清楚的例子,因為其效益非常具體:一個過去需要部署的時間變更,變成了一行的更新。但同樣的模式也適用於功能旗標、路由規則和速率限制。將設定儲存為資料,用變更串流監看它,並用恢復權杖處理重新連線。這三者的組合,是讓即時設定變得值得信賴,而不僅僅是方便的原因。