- 后端
- 工作流自动化
- 任务调度
【免费下载链接】temporal
Temporal service
导读
本文以 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.Invoke(api.go):先prepare校验请求并应用动态配置覆盖,然后prepareNewWorkflow创建新的 Mutable State;prepareNewWorkflow(api.go):通过api.NewWorkflowWithSignal构建初始 Mutable State,并调用mutableState.CloseTransactionAsSnapshot把「首个事件批次」落为快照。从源码看,它严格要求len(eventBatches) == 1,即首个批次恰好包含WorkflowExecutionStarted与WorkflowTaskScheduled两个事件,与文档所述一致;createBrandNew(api.go):以persistence.CreateWorkflowModeBrandNew模式调用CreateWorkflowExecution写入 executions 表,同时持久化 Mutable State 快照与 Transfer Task。
从源码还可看到一个值得注意的细节:prepare中会处理RequestEagerExecution(Eager 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 2:Worker 拉取并处理 Workflow Task
一个 Worker 拉取并处理 Workflow Task:
- 它推进工作流的执行,并在调用 Activity 处被阻塞(即工作流代码执行到
callActivity(myActivity)时暂停,等待 Activity 结果)。
对应状态视图:Mutable State 中 Workflow Task 变为Started,Transfer Queue 清空,Timer Queue 中新增Workflow Task Timeout,History 中追加第三个事件WorkflowTaskStarted:
源码入口:RecordWorkflowTaskStarted 的幂等与过期处理
Worker 通过 Matching 拿到任务后,由 Matching 调用 History 服务的RecordWorkflowTaskStarted。该 RPC 的 gRPC 入口在 service/history/handler.go(Handler.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 变为Scheduled;Transfer Queue 中出现Activity Task;Timer Queue 新增Activity Task Timeout;History 追加WorkflowTaskCompleted与ActivityTaskScheduled两个事件:
源码入口:RespondWorkflowTaskCompleted 与命令分发
Worker 以RespondWorkflowTaskCompleted请求把命令(Commands)带回 Frontend,最终进入 History 服务的WorkflowTaskCompletedHandler.Invoke(service/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 4:Worker 拉取并执行 Activity
一个 Worker 拉取 Activity Task 并执行该 Activity:
状态视图:Mutable State 中 Activity Task 变为Started,Transfer Queue 清空,Timer Queue 保留Activity Task Timeout与Workflow Timeout,History 追加ActivityTaskStarted(图中以虚线框标注该事件,示意其是可重放事件流中的普通一环):
源码入口:RecordActivityTaskStarted
Activity Task 的启动由 History 服务的RecordActivityTaskStartedRPC 记录,gRPC 入口在 service/history/handler.go。从 handler 源码可以看到一个当前仓库的重要实现细节:如果任务令牌(Task Token)带有组件引用(componentRef),该请求会被路由到 Chasm 引擎按独立 Activity(standalone activity)处理(handler.go);否则走常规的、由 Mutable State 支撑的工作流 Activity 路径,即engine.RecordActivityTaskStarted。
在常规路径下,History 服务追加ActivityTaskStarted事件、更新 Mutable State,并添加 Activity Task 超时定时任务(对应时序图中的 note 段)。Worker 随后真正执行业务逻辑(Execute activity),执行完成后准备回传结果。
六、Step 5:Activity 完成并调度下一个 Workflow Task
Activity 完成后,执行它的 Worker 发送RespondActivityTaskCompleted,其中包含 Activity 的结果:
- 一个新的 Workflow Task 被加入 Matching 服务。
状态视图:Mutable State 中 Workflow Task 回到Scheduled(新一轮工作流任务),Transfer Queue 中出现Workflow Task,Timer Queue 保留三个超时任务,History 追加ActivityTaskCompleted(携带 Activity 结果)与WorkflowTaskScheduled:
源码入口:RespondActivityTaskCompleted
gRPC 入口在 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 6:Worker 再次拉取 Workflow Task
Worker 拉取该 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.RespondWorkflowTaskCompleted(service/history/handler.go),内部逻辑与 Step 3 一致:WorkflowTaskCompletedHandler.Invoke校验令牌、追加WorkflowTaskCompleted事件,随后CompleteWorkflowExecution命令被处理,追加WorkflowExecutionCompleted事件,更新 Mutable State,并添加收尾类任务——包括更新可见性(Visibility)、上传分层存储归档(Tiered Storage,如上传到 S3)、执行保留期清理(Retention,删除数据)等。
QueueProcessor 最后一轮循环处理这些收尾任务,对应仓库中的三个独立执行器:
- 可见性更新:visibility_queue_task_executor.go(
ProcessTask (Update visibility)); - 归档上传:service/history/archival 目录及 archival_queue_task_executor.go(
Upload 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:
源码入口:RespondActivityTaskFailed
gRPC 入口在 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 服务(AddActivityTask),Worker 以 Attempt 2 重新执行该 Activity。这一「失败 → 追加事件 → 定时重试 → 重新调度」的循环会一直持续,直到 Activity 成功、达到最大重试次数(此时触发 Activity 最终失败,工作流按失败收尾),或工作流整体超时。
十、代码入口汇总与生命周期全景
将上文七个步骤的 RPC 与仓库源码入口汇总如下,便于按图索骥:
| 生命周期阶段 | RPC / 事件 | 源码入口(以仓库根目录为基准) |
|---|---|---|
| Step 1 启动工作流 | StartWorkflowExecution | service/history/api/startworkflow/api.go 中Starter.Invoke;队列投递见 service/history/queues/queue_immediate.go 与 transfer_queue_active_task_executor.go |
| Step 2 启动 Workflow Task | RecordWorkflowTaskStarted | service/history/handler.go;核心逻辑 service/history/api/recordworkflowtaskstarted/api.go |
| Step 3 调度 Activity | RespondWorkflowTaskCompleted+ScheduleActivityTask命令 | service/history/handler.go;命令处理 service/history/api/respondworkflowtaskcompleted/api.go |
| Step 4 启动 Activity | RecordActivityTaskStarted | service/history/handler.go |
| Step 5 Activity 完成 | RespondActivityTaskCompleted | service/history/handler.go |
| Step 6 再次拉取 Workflow Task | 同 Step 2 | 同 Step 2 |
| Step 7 完成工作流 | RespondWorkflowTaskCompleted+CompleteWorkflowExecution命令 | service/history/handler.go;收尾执行器:visibility_queue_task_executor.go、service/history/archival、deletemanager |
| 失败分支 | RespondActivityTaskFailed | service/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 5,Activity 完成,调度下一个 Workflow Task) → WorkflowTaskStarted → WorkflowTaskCompleted → WorkflowExecutionCompleted(Step 6–7,工作流收尾)贯穿始终的三条主线
从整个生命周期可以提炼出 Temporal 核心设计的三个不变规律:
- 一切状态变化都沉淀为事件:无论是任务调度、任务开始还是 Activity 完成,History 服务总是「先追加事件、再更新 Mutable State、再写任务」,任何一步都对应一次持久化事务,这是工作流可重放、可恢复的根本保障。
- Worker 与 History 通过「命令-事件」闭环协作:Worker 只返回命令(
ScheduleActivityTask、CompleteWorkflowExecution等),由 History 服务校验并转换为事件;Worker 自身不做任何持久化,因此任意 Worker 崩溃都不会破坏状态一致性。 - 队列是异步解耦的引擎:Transfer Queue 负责把任务搬运到 Matching,Timer 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 的事件构建逻辑,即可逐步打通整条调用链。
- 后端
- 工作流自动化
- 任务调度
【免费下载链接】temporal
Temporal service
相关推荐
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
后端前端企业应用
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考