
后端工作流自动化任务调度【免费下载链接】temporalTemporal service项目地址https://gitcode.com/gh_mirrors/te/temporal点击查看免费下载导读本文以 Temporal Server 官方架构文档 workflow-lifecycle.md 为主体完整还原一个最简单的「调用单个 Activity 并返回结果」的工作流在 Temporal Server 内部经历的全部阶段——从用户应用发出StartWorkflowExecution请求到 History 服务初始化历史、Matching 服务分发任务、Worker 执行并回传结果直至WorkflowExecutionCompleted事件落库。读完本文你将掌握 Workflow History、Mutable State、Workflow Task、Transfer Queue、Timer Queue 等核心概念的协作关系理解每一个关键 RPC 背后 History 服务触发的持久化写入并能结合仓库源码定位每个阶段的真实入口实现。术语约定下文凡是提到「初始化 History」「向 History 追加事件」或「持久化 Mutable State 与 History 任务」均指对持久化层的持久性写入durable write而非内存操作。一、先理解核心概念History、Mutable State 与两类 Task在进入分步流程前先厘清生命周期图中反复出现的几个概念它们是理解整个时序图的基础Workflow History工作流历史一条不可变的事件日志Event Log。工作流的每一次状态推进启动、任务调度、任务开始、任务完成、Activity 调度等都会以HistoryEvent的形式追加写入例如WorkflowExecutionStarted、WorkflowTaskScheduled、ActivityTaskScheduled、WorkflowExecutionCompleted。它也是 Worker 端事件溯源Event Sourcing重放工作流代码的依据。Mutable State可变状态工作流执行的最新内存态由 History 服务维护持久化到 executions 表中。它记录当前 Workflow Task / Activity Task 处于 Scheduled、Started 还是 Empty以及各类超时定时器等。Workflow Task / Activity Task任务调度单元。Workflow Task 由 Worker 拉取并执行工作流代码确定性重放Activity Task 由 Worker 拉取并执行真正的业务逻辑。Transfer Queue转移队列History 服务内部的即时任务队列immediate queue负责将「需要投递给 Matching 服务」的任务如 AddWorkflowTask、AddActivityTask异步搬运出去。对应源码 service/history/queues/queue_immediate.go。Timer Queue定时器队列存放各类超时任务Workflow Timeout、Workflow Task Timeout、Activity Task Timeout、Activity Retry由定时任务执行器处理。QueueProcessor队列处理器History 服务中轮询持久化任务表、取出任务并执行投递到 Matching、更新可见性、上传归档、删除数据等的后台循环。下文每个步骤中出现的loop QueueProcessor段代表 QueueProcessor 周期性地执行GetHistoryTasks → ProcessTask → AddWorkflowTask/AddActivityTask这一循环将历史任务表中的任务真正派发出去。二、Step 1启动工作流 —— StartWorkflowExecution用户应用发送StartWorkflowExecution请求这是整个生命周期的起点Workflow History 被初始化为两个事件[WorkflowExecutionStarted, WorkflowTaskScheduled]一个 Workflow Task 被加入 Matching 服务由 QueueProcessor 异步完成投递。此刻 History 服务与持久化层的状态视图如下Mutable State 中 Workflow Task 处于Scheduled状态Transfer Queue 中有一个WorkflowTask待投递Timer Queue 中挂起Workflow Timeout而 Workflow History 已写入前两个事件源码入口Starter.Invoke 与事务化创建前端Frontend服务校验请求后将StartWorkflowExecution转发给 History 服务最终进入 service/history/api/startworkflow/api.go 中的Starter结构体。其核心调用链如下Starter.Invokeapi.go先prepare校验请求并应用动态配置覆盖然后prepareNewWorkflow创建新的 Mutable StateprepareNewWorkflowapi.go通过api.NewWorkflowWithSignal构建初始 Mutable State并调用mutableState.CloseTransactionAsSnapshot把「首个事件批次」落为快照。从源码看它严格要求len(eventBatches) 1即首个批次恰好包含WorkflowExecutionStarted与WorkflowTaskScheduled两个事件与文档所述一致createBrandNewapi.go以persistence.CreateWorkflowModeBrandNew模式调用CreateWorkflowExecution写入 executions 表同时持久化 Mutable State 快照与 Transfer Task。从源码还可看到一个值得注意的细节prepare中会处理RequestEagerExecutionEager Workflow Start——若动态配置未开启或首个工作流任务存在 backoff则把 eager 标志置为 false回退到常规的 Matching 分发路径api.go。源码入口QueueProcessor 与 Transfer Queue持久化完成后[WorkflowExecutionStarted, WorkflowTaskScheduled]之外还写入了一条 Transfer Task。History 服务中的队列处理器循环GetHistoryTasks → ProcessTask → Matching.AddWorkflowTask由 service/history/queues/queue_immediate.go 中的immediateQueue实现——它以分页方式调用shard.GetHistoryTasks读取任务queue_immediate.go并通过独立的processEventLoop协程持续消费任务。实际的 Workflow Task 投递动作在 transfer_queue_active_task_executor.go 中完成。三、Step 2Worker 拉取并处理 Workflow Task一个 Worker 拉取并处理 Workflow Task它推进工作流的执行并在调用 Activity 处被阻塞即工作流代码执行到callActivity(myActivity)时暂停等待 Activity 结果。对应状态视图Mutable State 中 Workflow Task 变为StartedTransfer Queue 清空Timer Queue 中新增Workflow Task TimeoutHistory 中追加第三个事件WorkflowTaskStarted源码入口RecordWorkflowTaskStarted 的幂等与过期处理Worker 通过 Matching 拿到任务后由 Matching 调用 History 服务的RecordWorkflowTaskStarted。该 RPC 的 gRPC 入口在 service/history/handler.goHandler.RecordWorkflowTaskStarted核心实现在 service/history/api/recordworkflowtaskstarted/api.go 的Invoke函数中通过mutableState.GetWorkflowTaskByID(scheduledEventID)按 Scheduled Event ID 定位任务若任务已被其他请求完成会返回NotFound安全地丢弃本次任务这正是文档时序图中 Matching → History 一步背后的幂等保障api.go若任务已由同一 RequestID 启动过则直接复用结果返回updateAction.Noop true避免重复写入api.go若任务已被别的请求启动则返回TaskAlreadyStarted错误api.go。成功启动后Mutable State 事务会追加WorkflowTaskStarted事件并添加 Workflow Task 超时定时任务对应时序图中的 note 段同时 History 服务读取历史事件GetHistoryEvents随响应一起返回给 Worker供其重放工作流代码。Worker 端在执行工作流代码的过程中推进Advance workflow遇到callActivity时挂起等待下一步调度 Activity。四、Step 3调度 Activity —— ScheduleActivityTask 命令工作流发起 Activity 调用Worker 向 Frontend 回传一条ScheduleActivityTask命令一个 Activity Task 被加入 Matching 服务。状态视图Mutable State 中 Workflow Task 变为Empty当前工作流任务已处理完Activity Task 变为ScheduledTransfer Queue 中出现Activity TaskTimer Queue 新增Activity Task TimeoutHistory 追加WorkflowTaskCompleted与ActivityTaskScheduled两个事件源码入口RespondWorkflowTaskCompleted 与命令分发Worker 以RespondWorkflowTaskCompleted请求把命令Commands带回 Frontend最终进入 History 服务的WorkflowTaskCompletedHandler.Invokeservice/history/api/respondworkflowtaskcompleted/api.go。关键路径反序列化任务令牌Task Token通过一致性检查获取工作流租约GetWorkflowLeaseWithConsistencyCheck校验 Mutable State 中的 Workflow Task 与令牌信息ScheduledEventID、StartedEventID、StartedTime、Attempt、Version是否一致api.go调用ms.AddWorkflowTaskCompletedEvent追加WorkflowTaskCompleted事件api.go请求中的ScheduleActivityTask命令随后经命令处理分发最终在 Mutable State 中追加ActivityTaskScheduled事件、创建 Activity Task 的 Mutable State 条目并向持久化层写入新的 Transfer Task对应时序图中的 note 段。说明早期版本中ScheduleActivityTask命令的处理器位于service/history/workflow_task_handler.go当前仓库中该命令处理逻辑已重构进WorkflowTaskCompletedHandler命令分发链与 Mutable State 的事件构建逻辑中读者可在 service/history/api/respondworkflowtaskcompleted/api.go 与 service/history/workflow 目录下继续追踪。随后 QueueProcessor 再次循环将 Transfer Task 转化为Matching.AddActivityTask调用Activity Task 进入 Matching 服务等待 Worker 拉取。五、Step 4Worker 拉取并执行 Activity一个 Worker 拉取 Activity Task 并执行该 Activity状态视图Mutable State 中 Activity Task 变为StartedTransfer Queue 清空Timer Queue 保留Activity Task Timeout与Workflow TimeoutHistory 追加ActivityTaskStarted图中以虚线框标注该事件示意其是可重放事件流中的普通一环源码入口RecordActivityTaskStartedActivity Task 的启动由 History 服务的RecordActivityTaskStartedRPC 记录gRPC 入口在 service/history/handler.go。从 handler 源码可以看到一个当前仓库的重要实现细节如果任务令牌Task Token带有组件引用componentRef该请求会被路由到 Chasm 引擎按独立 Activitystandalone activity处理handler.go否则走常规的、由 Mutable State 支撑的工作流 Activity 路径即engine.RecordActivityTaskStarted。在常规路径下History 服务追加ActivityTaskStarted事件、更新 Mutable State并添加 Activity Task 超时定时任务对应时序图中的 note 段。Worker 随后真正执行业务逻辑Execute activity执行完成后准备回传结果。六、Step 5Activity 完成并调度下一个 Workflow TaskActivity 完成后执行它的 Worker 发送RespondActivityTaskCompleted其中包含 Activity 的结果一个新的 Workflow Task 被加入 Matching 服务。状态视图Mutable State 中 Workflow Task 回到Scheduled新一轮工作流任务Transfer Queue 中出现Workflow TaskTimer Queue 保留三个超时任务History 追加ActivityTaskCompleted携带 Activity 结果与WorkflowTaskScheduled源码入口RespondActivityTaskCompletedgRPC 入口在 service/history/handler.go。与启动 Activity 类似handler 会先反序列化任务令牌并校验若令牌含组件引用则交给 Chasm 处理否则走 Mutable State 支撑的常规路径engine.RespondActivityTaskCompleted。成功处理后持久化层追加ActivityTaskCompleted携带 Activity 结果与WorkflowTaskScheduled两个事件更新 Mutable State 并写入新的 Transfer Task对应时序图中的 note 段。QueueProcessor 再次循环将该 Transfer Task 转化为Matching.AddWorkflowTask新一轮 Workflow Task 进入 Matching 队列。七、Step 6Worker 再次拉取 Workflow TaskWorker 拉取该 Workflow Task它推进工作流发现工作流已经执行到末尾。这一步的时序与 Step 2 完全一致Same sequence diagram as step 2 above即复用「PollWorkflowTask → RecordWorkflowTaskStarted → 追加WorkflowTaskStarted事件 → 返回历史 → Worker 重放推进」的完整链路。区别在于本轮的 Mutable State 中不再有新的 Activity 需要调度工作流代码执行完毕后将返回最终结果进入收尾阶段。八、Step 7完成工作流 —— CompleteWorkflowExecution 命令Worker 发送RespondWorkflowTaskCompleted其中包含CompleteWorkflowExecution命令状态视图Mutable State 中 Workflow Task 变为Empty且不再有后续任务History 追加最终事件WorkflowExecutionCompleted完整事件序列见下图共 11 个事件源码入口收尾任务的落地本步的 gRPC 入口同样是Handler.RespondWorkflowTaskCompletedservice/history/handler.go内部逻辑与 Step 3 一致WorkflowTaskCompletedHandler.Invoke校验令牌、追加WorkflowTaskCompleted事件随后CompleteWorkflowExecution命令被处理追加WorkflowExecutionCompleted事件更新 Mutable State并添加收尾类任务——包括更新可见性Visibility、上传分层存储归档Tiered Storage如上传到 S3、执行保留期清理Retention删除数据等。QueueProcessor 最后一轮循环处理这些收尾任务对应仓库中的三个独立执行器可见性更新visibility_queue_task_executor.goProcessTask (Update visibility)归档上传service/history/archival 目录及 archival_queue_task_executor.goUpload to S3等分层存储动作数据清理由删除管理器与保留期清理逻辑完成Delete data相关代码见 deletemanager。九、分支场景Activity 失败与重试上述 Step 1–7 描述了最顺利的执行路径。当 Activity 执行失败时流程走向另一条分支Activity 可能失败并被重试状态视图Mutable State 中 Activity Task 变为Scheduled, Attempt 2第二次尝试Timer Queue 中出现Activity Retry定时任务History 追加ActivityTaskFailed与ActivityTaskScheduled源码入口RespondActivityTaskFailedgRPC 入口在 service/history/handler.go。与完成路径对称反序列化令牌、按组件引用分流到 Chasm 或常规路径最终由engine.RespondActivityTaskFailed处理。持久化层追加ActivityTaskFailed与ActivityTaskScheduled两个事件更新 Mutable State并添加Activity Retry 定时任务对应时序图 note 段中的add Timer Task (activity timout)。此后当重试定时器到期由 timer_queue_active_task_executor.go 处理QueueProcessor 会把新的 Activity Task 再次投递到 Matching 服务AddActivityTaskWorker 以 Attempt 2 重新执行该 Activity。这一「失败 → 追加事件 → 定时重试 → 重新调度」的循环会一直持续直到 Activity 成功、达到最大重试次数此时触发 Activity 最终失败工作流按失败收尾或工作流整体超时。十、代码入口汇总与生命周期全景将上文七个步骤的 RPC 与仓库源码入口汇总如下便于按图索骥生命周期阶段RPC / 事件源码入口以仓库根目录为基准Step 1 启动工作流StartWorkflowExecutionservice/history/api/startworkflow/api.go 中Starter.Invoke队列投递见 service/history/queues/queue_immediate.go 与 transfer_queue_active_task_executor.goStep 2 启动 Workflow TaskRecordWorkflowTaskStartedservice/history/handler.go核心逻辑 service/history/api/recordworkflowtaskstarted/api.goStep 3 调度 ActivityRespondWorkflowTaskCompletedScheduleActivityTask命令service/history/handler.go命令处理 service/history/api/respondworkflowtaskcompleted/api.goStep 4 启动 ActivityRecordActivityTaskStartedservice/history/handler.goStep 5 Activity 完成RespondActivityTaskCompletedservice/history/handler.goStep 6 再次拉取 Workflow Task同 Step 2同 Step 2Step 7 完成工作流RespondWorkflowTaskCompletedCompleteWorkflowExecution命令service/history/handler.go收尾执行器visibility_queue_task_executor.go、service/history/archival、deletemanager失败分支RespondActivityTaskFailedservice/history/handler.go重试定时器 timer_queue_active_task_executor.go以本文的示例工作流调用单个 Activity 并返回结果为例完整成功路径的 History 事件序列为WorkflowExecutionStarted → WorkflowTaskScheduled → WorkflowTaskStarted → WorkflowTaskCompleted Step 1–2首个 Workflow Task 执行工作流代码 → ActivityTaskScheduled → ActivityTaskStarted Step 3–4调度并执行 Activity → ActivityTaskCompleted → WorkflowTaskScheduled Step 5Activity 完成调度下一个 Workflow Task → WorkflowTaskStarted → WorkflowTaskCompleted → WorkflowExecutionCompletedStep 6–7工作流收尾贯穿始终的三条主线从整个生命周期可以提炼出 Temporal 核心设计的三个不变规律一切状态变化都沉淀为事件无论是任务调度、任务开始还是 Activity 完成History 服务总是「先追加事件、再更新 Mutable State、再写任务」任何一步都对应一次持久化事务这是工作流可重放、可恢复的根本保障。Worker 与 History 通过「命令-事件」闭环协作Worker 只返回命令ScheduleActivityTask、CompleteWorkflowExecution等由 History 服务校验并转换为事件Worker 自身不做任何持久化因此任意 Worker 崩溃都不会破坏状态一致性。队列是异步解耦的引擎Transfer Queue 负责把任务搬运到 MatchingTimer Queue 负责超时与重试Visibility/Archival 队列负责收尾。所有投递与副作用都以持久化任务为媒介异步完成保证了系统的高可用与可扩展性。结语本文从 docs/architecture/workflow-lifecycle.md 出发完整走读了一个最小工作流从启动到完成的七个阶段并补充了失败重试分支。每个阶段的时序图、状态视图与仓库源码入口service/history/handler.go、service/history/api 下的各 API 实现、service/history/queues 队列实现共同构成了理解 Temporal Server 内部运转的完整地图。对于想深入源码的读者建议从Starter.Invoke与WorkflowTaskCompletedHandler.Invoke两个入口入手配合 service/history/workflow 目录下 Mutable State 的事件构建逻辑即可逐步打通整条调用链。赞分享后端工作流自动化任务调度【免费下载链接】temporalTemporal service项目地址https://gitcode.com/gh_mirrors/te/temporal点击查看免费下载相关推荐Thor机械臂3D打印攻略快速掌握STL文件使用与打印技巧Thor机械臂3D打印攻略快速掌握STL文件使用与打印技巧 Thor是一款开源3D打印机械臂拥有6个自由度高度达625mm最大负载750g非常适合教育硬件开发机器人智能硬件OpenMetadata 事件生命周期工作流将数据质量事故管理从硬编码状态机迁移到 Flowable 治理工作流OpenMetadata 事件生命周期工作流将数据质量事故管理从硬编码状态机迁移到 Flowable 治理工作流 本文基于 OpenMetadata 的变更提数据目录数据血缘数据治理后端MCP 服务Cal.diy Embed 生命周期详解握手协议、命令队列与事件状态机Cal.diy Embed 生命周期详解握手协议、命令队列与事件状态机 导读 Cal.diy Scheduling infrastructure for a后端前端企业应用上一篇给网页加粒子背景只要10分钟tsParticles 粒子动画库实战指南下一篇终极指南如何在Seal中配置yt-dlp自定义提取器插件创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考