行語義(Durable Execution):分布式工作流的每一步狀態(tài)都可靠落盤)
Conductor 持久化執(zhí)行語義Durable Execution分布式工作流的每一步狀態(tài)都可靠落盤【免費下載鏈接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents項目地址: https://gitcode.com/GitHub_Trending/co/conductor本文是 Conductor一個面向應(yīng)用與 AI Agent 的事件驅(qū)動工作流引擎的持久化執(zhí)行Durable Execution語義權(quán)威指南。文章圍繞每個工作流執(zhí)行在每個步驟都會被持久化、能夠抵御基礎(chǔ)設(shè)施故障、并保證任務(wù)至少一次at-least-once投遞這一核心模型展開完整覆蓋持久化內(nèi)容清單、任務(wù)投遞保證、故障矩陣、任務(wù)狀態(tài)機、超時與重試配置、工作流級耐久性、重放與恢復(fù)以及分布式一致性。讀完本文你將掌握 Conductor 如何在服務(wù)器重啟、Worker 崩潰、網(wǎng)絡(luò)分區(qū)等各類故障下不丟失執(zhí)行進度以及如何據(jù)此設(shè)計冪等的 Worker 與可安全升級的工作流定義。什么是 Durable Execution引擎的可靠性底座Conductor 本質(zhì)上是一個面向分布式工作流與持久化 Agent 的執(zhí)行引擎。所謂 Durable Execution持久化執(zhí)行指的是工作流的每一次執(zhí)行都在每一步被完整持久化執(zhí)行進度不會因進程崩潰、節(jié)點重啟或機房故障而丟失。結(jié)合任務(wù)隊列、重試與超時機制引擎能夠保證任務(wù)至少一次投遞從而讓工作流與 Agent 永不丟失進度成為可驗證的工程承諾。從代碼結(jié)構(gòu)看這一模型貫穿引擎核心TaskModel.java 定義了任務(wù)的數(shù)據(jù)模型與狀態(tài)機Status枚舉每個任務(wù)實例攜帶scheduledTime、startTime、endTime、updateTime、retryCount、pollCount等時間戳與計數(shù)正是這些字段支撐了后續(xù)的持久化、超時判定與重試WorkflowModel.java 定義了工作流實例的狀態(tài)模型DeciderService.java 與 WorkflowSweeper.java 等實現(xiàn)了decide 求值 sweeper 清掃的推進機制。每一步都持久化What Persists當(dāng)一個工作流執(zhí)行時Conductor 會持久化以下四類數(shù)據(jù)工作流定義快照Workflow definition snapshot——本次執(zhí)行所使用的定義副本啟動后即不可變immutable工作流狀態(tài)Workflow state——狀態(tài)、輸入、輸出、correlation ID 以及變量variables每一次任務(wù)執(zhí)行Every task execution——狀態(tài)、輸入、輸出、時間戳、重試次數(shù)與 worker ID任務(wù)隊列狀態(tài)Task queue state——哪些任務(wù)處于已調(diào)度、進行中或已完成。全部狀態(tài)都會在進入下一步之前寫入所配置的持久化存儲Redis、PostgreSQL、MySQL 或 Cassandra。如果服務(wù)器重啟執(zhí)行將從最后持久化的狀態(tài)恢復(fù)而不是從頭開始。定義快照這一設(shè)計意義重大它意味著運行中的執(zhí)行與元數(shù)據(jù)存儲metadata store解耦。即使工作流定義在運行期間被更新甚至被刪除正在運行的實例依然使用自己內(nèi)嵌的快照繼續(xù)執(zhí)行——這直接支撐了零停機升級。任務(wù)投遞保證At-Least-Once DeliveryConductor 對所有任務(wù)提供至少一次投遞保證其循環(huán)如下任務(wù)被調(diào)度時放入持久化任務(wù)隊列persistent task queueWorker 輪詢poll并領(lǐng)取任務(wù)任務(wù)進入IN_PROGRESSWorker 完成任務(wù)后上報COMPLETEDConductor 推進工作流若 Worker 失敗或崩潰任務(wù)將基于重試與超時配置被重新投遞redelivered。一個任務(wù)永遠不會被靜默丟失。如果 Worker 領(lǐng)取了任務(wù)卻始終不響應(yīng)**響應(yīng)超時response timeout**會觸發(fā)重新投遞。值得注意的是這同時是至少一次而非恰好一次——同一任務(wù)可能被執(zhí)行多次。因此Worker 必須以冪等的方式處理副作用這一點在文末這對你的代碼意味著什么一節(jié)有進一步說明。故障矩陣每種故障場景下引擎的確切行為下表給出了 Conductor 在各類故障場景下的精確行為場景Conductor 的行為結(jié)果Worker 輪詢后在開始任何工作前崩潰觸發(fā)響應(yīng)超時response timeout。任務(wù)回到SCHEDULED由新 Worker 領(lǐng)取。任務(wù)自動重試無數(shù)據(jù)丟失。Worker 在產(chǎn)生副作用之后、上報完成之前崩潰觸發(fā)響應(yīng)超時。任務(wù)被重新投遞給另一個 Worker。任務(wù)會再次執(zhí)行。Worker 必須對副作用冪等或使用任務(wù)的updateTime檢測重投遞。Worker 上報 FAILEDConductor 根據(jù)重試配置retryCount、retryDelaySeconds、retryLogic創(chuàng)建一次新的任務(wù)執(zhí)行。重試至配置上限。重試耗盡后任務(wù)進入FAILED工作流的失敗處理邏輯接管。Worker 上報 FAILED_WITH_TERMINAL_ERROR不重試任務(wù)立即終止。工作流失敗或執(zhí)行配置的failureWorkflow。工作流執(zhí)行期間服務(wù)器重啟重啟后sweeper 服務(wù)從持久化存儲拾取進行中的工作流并重新求值re-evaluate。從最后持久化的狀態(tài)恢復(fù)執(zhí)行無需人工干預(yù)??缍啻尾渴鸬拈L時間等待WAIT 與 HUMAN 任務(wù)在持久化存儲中保持IN_PROGRESS。計時器或信號的解析是持久的。當(dāng)?shù)却龝r長耗盡或信號到達時即使是在多次部署之后的數(shù)天任務(wù)完成工作流繼續(xù)推進。暫停中的工作流收到信號/WebhookTask Update API 或事件處理器將 WAIT/HUMAN 任務(wù)置為COMPLETED并提供輸出。工作流立即恢復(fù)信號載荷可作為任務(wù)輸出使用。運行期間工作流定義被更新運行中的執(zhí)行繼續(xù)使用啟動時拍攝的定義快照新執(zhí)行使用新定義。定義變更不影響任何運行中的執(zhí)行實現(xiàn)零停機升級。運行期間工作流版本被刪除運行中的執(zhí)行與元數(shù)據(jù)存儲解耦繼續(xù)使用內(nèi)嵌的定義快照?,F(xiàn)有執(zhí)行正常完成只有新啟動的實例受影響。Worker 與服務(wù)器之間網(wǎng)絡(luò)分區(qū)Worker 的更新無法到達服務(wù)器觸發(fā)響應(yīng)超時任務(wù)重新入隊。分區(qū)恢復(fù)后新 Worker或同一個 Worker重新領(lǐng)取任務(wù)。這張矩陣是本文檔的靈魂它用一張表窮舉了任何環(huán)節(jié)出問題都不會丟進度的承諾邊界要么自動重試要么轉(zhuǎn)交失敗處理要么等待恢復(fù)——但絕不會出現(xiàn)任務(wù)憑空消失的狀態(tài)。任務(wù)狀態(tài)機Task State Transitions每個任務(wù)遵循以下狀態(tài)機SCHEDULED ──→ IN_PROGRESS ──→ COMPLETED │ │ │ ├──→ FAILED ──→ SCHEDULED (retry) │ │ │ ├──→ FAILED_WITH_TERMINAL_ERROR │ │ │ └──→ TIMED_OUT ──→ SCHEDULED (retry) │ └──→ CANCELED (workflow terminated)**終態(tài)Terminal states**包括COMPLETED、FAILED重試耗盡后、FAILED_WITH_TERMINAL_ERROR、CANCELED、COMPLETED_WITH_ERRORS可選任務(wù) optional tasks。每一次狀態(tài)轉(zhuǎn)移都會在任何后續(xù)動作發(fā)生之前被持久化。在源碼中這套狀態(tài)機被建模為 TaskModel.Status 枚舉并且每個狀態(tài)顯式標(biāo)注了三個關(guān)鍵屬性狀態(tài)terminalsuccessfulretriableIN_PROGRESS否是是CANCELED是否否FAILED是否是FAILED_WITH_TERMINAL_ERROR是否否COMPLETED是是是COMPLETED_WITH_ERRORS是是是SCHEDULED否是是TIMED_OUT是否是SKIPPED是是否這一枚舉定義直接印證了文檔中的行為描述FAILED_WITH_TERMINAL_ERROR的retriablefalse因此引擎不會對它重試FAILED與TIMED_OUT的retriabletrue因此會按配置重試后重新回到SCHEDULED。isRetriable()、isSuccessful()、isTerminal()這三個方法就是引擎在 decide 流程中判斷是否該重試、是否算成功、是否已到終態(tài)的依據(jù)。超時與重試配置每任務(wù)級參數(shù)耐久性可以通過任務(wù)定義task definition按任務(wù)單獨配置詳見 taskdef.md。核心參數(shù)如下參數(shù)作用timeoutSeconds任務(wù)到達終態(tài)所允許的最大墻鐘時間wall-clock time。responseTimeoutSeconds在重新入隊前等待 Worker 狀態(tài)更新的最大時間。pollTimeoutSeconds一個已調(diào)度任務(wù)在被輪詢前等待的最大時間超時即觸發(fā)超時。retryCount失敗或超時時的重試次數(shù)。retryLogicFIXED、EXPONENTIAL_BACKOFF或LINEAR_BACKOFF。retryDelaySeconds重試之間的基礎(chǔ)延遲。timeoutPolicyRETRY、TIME_OUT_WF或ALERT_ONLY。從源碼 TaskDef.java 可以看到這些枚舉與默認(rèn)值的真實定義public enum TimeoutPolicy { RETRY, TIME_OUT_WF, ALERT_ONLY } public enum RetryLogic { FIXED, EXPONENTIAL_BACKOFF, LINEAR_BACKOFF }其默認(rèn)值分別為retryCount默認(rèn)為3timeoutPolicy默認(rèn)為TIME_OUT_WF即任務(wù)超時后直接判定工作流超時失敗retryLogic默認(rèn)為FIXEDretryDelaySeconds默認(rèn)為60秒timeoutSeconds無默認(rèn)值需顯式配置且?guī)otNull校驗約束。理解這幾個默認(rèn)值有助于避免我以為不會重試、結(jié)果重試了 3 次或我以為會重試、結(jié)果工作流直接超時失敗之類的配置誤區(qū)。其中retryDelaySeconds是三種重試邏輯共用的基礎(chǔ)延遲FIXED每次固定等待該時長EXPONENTIAL_BACKOFF按指數(shù)遞增LINEAR_BACKOFF按線性遞增。responseTimeoutSeconds與pollTimeoutSeconds則共同決定了多久判定一個 Worker 失聯(lián)、多久判定一個任務(wù)無人領(lǐng)取是故障矩陣中Worker 崩潰后自動重投遞得以實現(xiàn)的計時器基礎(chǔ)。工作流級耐久性超越單任務(wù)除單個任務(wù)外Conductor 還提供工作流級別的耐久能力補償流Compensation flows配置一個failureWorkflow當(dāng)主工作流失敗時自動運行并攜帶完整上下文失敗原因、失敗任務(wù) ID、工作流執(zhí)行數(shù)據(jù)暫停與恢復(fù)Pause and resume任意運行中的工作流可通過 API 暫停并在之后恢復(fù)狀態(tài)被完整保留重啟、重跑與重試Restart / rerun / retry詳見下文重放與恢復(fù)一節(jié)版本化Versioning多個工作流版本可以并發(fā)運行運行中的執(zhí)行對定義變更不可變重啟時可選地使用最新定義。這里的failureWorkflow是實現(xiàn)Saga 補償模式的官方入口主流程失敗后自動觸發(fā)補償流程撤銷已完成的副作用而補償流程本身同樣享受整套持久化保證。重放與恢復(fù)Replay and Recovery每一個工作流執(zhí)行都是完全可重放的fully replayable。Conductor 保留了完整的執(zhí)行圖——每個任務(wù)的輸入、輸出與狀態(tài)——因此你可以隨時重新執(zhí)行工作流。操作作用適用場景Restart重啟從開頭重新執(zhí)行整個工作流定義已變更需要一次干凈的執(zhí)行Rerun重跑從某個特定任務(wù)開始重新執(zhí)行復(fù)用之前任務(wù)的輸出修復(fù)中間某個任務(wù)而無需重跑全部Retry重試重試最后一個失敗的任務(wù)并從該點繼續(xù)瞬時故障、外部依賴當(dāng)時不可用這三個操作都可以作用于任意終態(tài)COMPLETED、FAILED、TIMED_OUT、TERMINATED的工作流并且可以無限期使用——因為 Conductor 完整保留了執(zhí)行圖。Restart 還可以選擇性地使用最新的工作流定義這樣你可以在修復(fù)定義中的 bug 后立刻重放執(zhí)行。分布式一致性多節(jié)點部署下的正確性在多節(jié)點部署中Conductor 通過以下機制保證一致性分布式鎖Distributed locking在整個集群中每個工作流同一時刻只有一個decide求值在運行可插拔實現(xiàn)Zookeeper、Redis柵欄令牌Fencing tokens防止持有過期鎖的節(jié)點提交過期更新持久化隊列Persistent queues任務(wù)隊列在節(jié)點故障后依然存活。支持可配置的分片策略round-robin 或 local-only在分布性與一致性之間做權(quán)衡。分布式鎖配置詳見部署指南中的 locking 一節(jié)。在源碼層面這對應(yīng) WorkflowReconciler.java、WorkflowSweeper.java 與 ExecutionLockService.java 等組件sweeper 定期從存儲中拾取未決工作流在分布式鎖的保護下執(zhí)行 decide確保同一工作流的推進在任何時刻只發(fā)生在一個節(jié)點上從而避免雙份推進導(dǎo)致的狀態(tài)錯亂。這對你的代碼意味著什么Worker 應(yīng)該是冪等的。由于至少一次投遞保證任務(wù)可能被執(zhí)行不止一次。請將 Worker 設(shè)計為能夠安全處理重投遞。你不需要自己構(gòu)建重試邏輯。Conductor 負(fù)責(zé)重試、超時與重新入隊。你的 Worker 只需上報成功或失敗。長時間運行的流程是安全的。使用 WAIT 與 HUMAN 任務(wù)來處理跨越數(shù)分鐘到數(shù)天的暫停狀態(tài)在多次部署之間保持持久。定義變更是安全的。可以隨時更新工作流定義而不影響正在運行的執(zhí)行。以零停機的方式逐步發(fā)布新版本。作為補充對于Worker 在副作用之后崩潰的場景文檔給出了兩條工程路徑一是讓 Worker 對副作用冪等重復(fù)執(zhí)行同一副作用的結(jié)果相同二是利用任務(wù)上的updateTime字段見 TaskModel.java 中的字段定義檢測這個任務(wù)是否已經(jīng)處理過一次從而在重投遞時跳過已完成的副作用。這一字段隨任務(wù)狀態(tài)一同持久化正是 Durable Execution 模型中可檢測的重投遞與業(yè)務(wù)冪等之間的銜接點。總而言之Conductor 的持久化執(zhí)行語義可以濃縮為一句話每一步都落盤任何故障都有明確的、可預(yù)期的行為路徑?;谶@一模型你可以放心地把跨機器、跨進程、跨部署的工作流編排任務(wù)交給引擎而把精力集中在業(yè)務(wù)邏輯本身?!久赓M下載鏈接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents項目地址: https://gitcode.com/GitHub_Trending/co/conductor創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考