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