引言
传统写法里,一个「下单后等支付、支付后发货、超时未支付则取消」的流程会被拆成若干个回调接口加一张状态表,再加一个定时任务扫超时。这套东西能跑,但每加一个步骤,状态分支就翻一倍,最终没人能说清「订单处于状态 X 时收到事件 Y 应该怎么办」。
Temporal 提出的解法叫持久化执行(Durable Execution):你写的就是一段普通的、同步的、带 if-else 和 try-catch 的代码,引擎把这段代码的每一次副作用记录成事件,进程崩溃后从事件历史重放,把执行恢复到崩溃前的状态。于是「等待 30 天」在代码里就是 Workflow.sleep(30 天),「失败重试」就是给 Activity 配一个 RetryPolicy,不需要状态表,不需要回调接口。
代价是引入了一个反直觉的约束:Workflow 代码必须是确定性的。同一段代码在重放时必须走完全相同的分支,否则状态会漂移。这个约束是 Temporal 所有「坑」的根源,也是理解它的关键。
本文先讲清楚重放模型到底是怎么工作的,再逐条拆解确定性约束,然后覆盖 Workflow 与 Activity 的职责划分、任务队列、信号、定时器、重试、Saga、版本兼容与历史截断,最后给出部署运维与落地路线。想先看整体选型框架的读者,可以从 工作流引擎全景与选型 开始。
目录
- 持久化执行要解决什么问题
- 事件溯源与重放模型
- 确定性约束清单
- Workflow 与 Activity 的职责划分
- 一个完整的订单工作流
- Worker 与任务队列
- 启动、查询与信号
- 定时器与长睡眠
- 重试策略与超时
- Saga 补偿模式
- 子工作流与并行执行
- 版本兼容与 GetVersion
- ContinueAsNew 与历史截断
- 本地活动与优化手段
- 数据转换与序列化
- 与消息中间件的集成
- 与外部系统的幂等桥接
- 部署与运维
- 自建与 Temporal Cloud 的取舍
- 观测与调试
- 落地路线图
- 权衡取舍
- 常见坑清单
- 小结
1. 持久化执行要解决什么问题
一个跨天的业务流程,用传统写法要处理五类问题:进程崩溃后状态从哪恢复、外部回调如何找到正确的实例、超时如何扫描与触发、重试如何保证不重复执行、流程逻辑变更后老实例怎么办。这五类问题与业务逻辑无关,但每个项目都要重写一遍。
持久化执行的思路是:把「执行进度」本身变成一等公民。引擎在执行过程中不断把「刚才做了什么」写进持久化的事件历史,进程崩溃后不从头开始,而是读历史、跳过已完成的部分、从断点继续。这样上面五类问题全部由引擎承担。
这个模型有一个必要前提:重放要能重建出相同的状态。所以代码不能依赖「当前时间」「随机数」「外部服务返回值」这些在重放时不一致的东西——它们必须先被记录进事件历史,重放时从历史里读。这就是确定性约束的由来。
执行过程(真实): Activity(扣库存) -> Activity(扣款) -> Timer(30天) -> ...
事件历史(持久): ActivityTaskScheduled(扣库存)
ActivityTaskCompleted(结果=OK)
ActivityTaskScheduled(扣款)
ActivityTaskCompleted(结果=OK)
TimerStarted(30天)
崩溃 -> 重放:按历史回放,遇到未完成的事件才真正执行
2. 事件溯源与重放模型
Temporal 的状态存储是「事件历史 + 可变状态」的组合。事件历史(Event History)是不可变的追加日志,记录了 Workflow 执行过程中的所有决策与结果;可变状态(Mutable State)是引擎为加速重放而维护的缓存,可以从事件历史重建。
每次 Workflow 代码需要「决策」(下一步调哪个 Activity、开哪个 Timer)时,SDK 会先检查历史里有没有对应的命令。如果有,就用历史里的结果继续;如果没有,才把命令发给服务端并等待结果。这个过程对开发者完全透明,你看到的只是一段同步代码。
// 开发者视角:一段同步代码
activities.charge(orderId);
// 引擎视角:第一次执行时发送 ScheduleActivityTask 命令并等待;
// 重放时直接从历史读取已完成结果,不产生副作用
重放的触发时机有四种:Worker 崩溃重启、Workflow 被缓存淘汰后重新加载、接收到新信号、定时器触发。前两种是「纯重放」(不执行新逻辑),后两种是「重放到断点后继续执行」。
理解这一点很重要:重放会重复执行 Workflow 代码,但不会重复执行 Activity。所以 Workflow 代码里的所有副作用(打印日志、写数据库、调 HTTP)都是危险的,它们会在每次重放时再执行一遍。
3. 确定性约束清单
Workflow 代码里不能做的事,以及正确的替代方式:
| 禁止 | 原因 | 替代 |
|---|---|---|
System.currentTimeMillis() | 重放时时间不同,分支漂移 | Workflow.currentTimeMillis() |
new Random() / Math.random() | 重放时结果不同 | Workflow.randomUUID() |
| 直接调 HTTP / 数据库 | 重放会重复副作用 | 放进 Activity |
| 直接读文件 / 环境变量 | 环境可能已变 | 放进 Activity 或作为输入参数 |
Thread.sleep() | 阻塞 Worker 线程 | Workflow.sleep() |
遍历 HashMap 做有序决策 | 迭代顺序不稳定 | 用 TreeMap 或先排序 |
| 多线程 / 线程池 | 重放顺序不可控 | Async.function / Promise |
| 依赖外部可变状态做分支 | 重放时值已变 | 把值存进 Workflow 变量 |
最容易踩的是「读环境变量」和「遍历 Map」。前者在本地和线上环境不同时表现不一致;后者在 Java 里对 HashMap 的迭代顺序在插入不同 key 时可能不同,导致两个 Worker 重放出不同的分支,触发 NonDeterminismError。
还有一类隐蔽的违反:在 Workflow 里调用一个「看起来是纯函数」的工具类,但这个工具类内部读了下游服务的配置。这类问题在代码评审时几乎看不出来,只能靠「Workflow 类里只允许出现 SDK 提供的 API 与 Activity 接口」这条硬规则来防。
4. Workflow 与 Activity 的职责划分
划分原则很简单:任何与外部世界交互的动作都放进 Activity,Workflow 只做决策。
- Workflow 负责:控制流(循环、分支、并行)、状态变量、定时器、信号处理、调用 Activity。
- Activity 负责:调外部 API、读写数据库、发消息、文件操作、任何可能失败且需要重试的副作用。
Activity 必须假定自己是「至少执行一次」的。引擎在崩溃恢复时可能重复投递同一个 Activity,所以每个 Activity 都要幂等,通常用一个业务幂等键(业务单号 + 步骤名)去重,细节见 重试幂等与补偿设计
。
Activity 还要注意超时设置。Temporal 有四类超时:ScheduleToClose(从调度到完成的总时限)、ScheduleToStart(在队列里等待的时限)、StartToClose(单次执行的时限)、Heartbeat(心跳间隔)。默认的 StartToClose 是 10 分钟,一个跑 40 分钟的批处理会被强制超时,必须显式配置。
@ActivityInterface
public interface OrderActivities {
@ActivityMethod(scheduleToCloseTimeout = Duration.ofMinutes(5),
startToCloseTimeout = Duration.ofMinutes(2),
retryPolicy = @RetryPolicy(maximumAttempts = 3,
initialInterval = @Interval(seconds = 1),
backoffCoefficient = 2.0))
void charge(String orderId, BigDecimal amount);
}
5. 一个完整的订单工作流
把前面几节的概念拼成一个可运行的例子:
public class OrderWorkflowImpl implements OrderWorkflow {
private final OrderActivities activities =
Workflow.newActivityStub(OrderActivities.class, ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofMinutes(2))
.setRetryPolicy(RetryOptions.newBuilder()
.setInitialInterval(Duration.ofSeconds(1))
.setBackoffCoefficient(2.0)
.setMaximumAttempts(5)
.build())
.build());
private boolean shipped = false;
@Override
public void execute(OrderInput input) {
Saga saga = new Saga(new Saga.Options.Builder().build());
saga.addCompensation(activities::unlockStock, input.orderId());
activities.lockStock(input.orderId());
saga.addCompensation(activities::refund, input.orderId());
activities.charge(input.orderId(), input.amount());
// 最多等 30 天,收到 signal 提前唤醒
boolean ok = Workflow.await(Duration.ofDays(30), () -> this.shipped);
if (!ok) {
saga.compensate();
return;
}
activities.notifyUser(input.orderId());
}
@Override
public void markShipped() {
this.shipped = true;
}
}
这段代码里有三个关键点:Workflow.newActivityStub 在 Workflow 构造时创建(不是每步创建,因为 Stub 本身是确定性的);Saga 的补偿按注册的逆序执行;Workflow.await 的条件里读的是 Workflow 实例变量,这个变量会被持久化进事件历史,所以重放时能恢复。
6. Worker 与任务队列
Temporal 的 Worker 是「拉取任务并执行」的进程。它轮询一个或多个任务队列(Task Queue),拿到 Workflow Task 时执行 Workflow 代码,拿到 Activity Task 时执行 Activity 代码。
WorkerFactory factory = WorkerFactory.newInstance(client);
Worker worker = factory.newWorker("order-task-queue");
worker.registerWorkflowImplementationTypes(OrderWorkflowImpl.class);
worker.registerActivitiesImplementations(new OrderActivitiesImpl());
factory.start();
任务队列的设计原则:
- 按业务域划分(
order-task-queue、payment-task-queue),不要所有流程共用一个。 - 按资源需求划分:需要 GPU 的 Activity 单独一个队列,跑在专门的 Worker 池上。
- Activity 队列可以独立伸缩,Workflow 队列的 Worker 数量可以很少(因为 Workflow 代码执行很快,大部分时间在等)。
一个常见的性能问题:把耗时 Activity 和快速 Activity 放在同一个队列,慢任务占满 Worker 的并发槽位,快任务排队等待。解决方式是拆队列,而不是加大并发数,因为并发数受 Worker 内存限制。
7. 启动、查询与信号
客户端启动 Workflow 时,WorkflowId 是幂等的关键:
WorkflowOptions options = WorkflowOptions.newBuilder()
.setWorkflowId("order-" + orderId) // 业务幂等键
.setWorkflowIdReusePolicy(
WorkflowIdReusePolicy.WORKFLOW_ID_REUSE_POLICY_REJECT_DUPLICATE)
.setTaskQueue("order-task-queue")
.build();
WorkflowStub stub = client.newUntypedWorkflowStub("OrderWorkflow", options);
stub.start(input);
WorkflowId 在默认配置下全局唯一,重复启动会返回已存在的实例(或按 ReusePolicy 拒绝),这是最自然的幂等启动方式。不要用随机 UUID 做 WorkflowId,那样失去幂等能力。
信号(Signal)是外部向 Workflow 推送事件的方式,比如支付回调通知「已支付」、仓库通知「已发货」。查询(Query)是只读地读取 Workflow 当前状态,不会改变历史,适合做「查一下这个订单现在走到哪了」。
# CLI:发信号、查状态
temporal workflow signal --workflow-id order-1001 --name markShipped
temporal workflow query --workflow-id order-1001 --name getStatus
temporal workflow describe --workflow-id order-1001
8. 定时器与长睡眠
Workflow.sleep(duration) 是持久化定时器,不是 Thread.sleep。它会写一条 TimerStarted 事件,Worker 立刻释放线程,时间到了再由服务端投递新的 Workflow Task。
Workflow.sleep(Duration.ofDays(7)); // 睡 7 天
boolean done = Workflow.await(Duration.ofHours(2), () -> ready); // 条件等待带超时
一个反直觉的性能细节:Workflow.sleep 的精度受 Workflow Task 超时与 Worker 缓存影响,通常在毫秒到秒级,不适合做「精确到毫秒」的调度。如果需要精确定时,应该用 Timer 唤醒后再判断,或者干脆用专门的任务调度器。
另一个细节是「睡很久的 Workflow 会占用缓存」。Temporal 的 Worker 会缓存一定数量的 Workflow 实例,睡 30 天的实例如果一直在缓存里会挤占活跃实例。引擎会在缓存压力大时把空闲实例从内存卸载(只是卸载,状态仍在服务端),唤醒时再重放。
9. 重试策略与超时
Temporal 的重试分两层:Activity 层的自动重试(由 RetryPolicy 控制),以及 Workflow 层的业务重试(由代码控制)。
RetryOptions retry = RetryOptions.newBuilder()
.setInitialInterval(Duration.ofSeconds(1))
.setBackoffCoefficient(2.0)
.setMaximumInterval(Duration.ofMinutes(1))
.setMaximumAttempts(10)
.setNonRetryableErrorTypes(List.of("InvalidArgumentError"))
.build();
setNonRetryableErrorTypes 非常关键:业务性错误(参数非法、余额不足)不该重试,只有技术性错误(网络超时、下游 5xx)才重试。不区分这两类错误,会导致无意义的重试风暴。
重试的终止条件要设计好。MaximumAttempts 用完之后,Activity 会抛出 ActivityFailure,Workflow 代码要捕获它并决定下一步:是走补偿、是转人工处理、还是标记为失败终态。默认不捕获会让整个 Workflow 失败,这通常不是业务想要的。
10. Saga 补偿模式
Temporal 的 Java SDK 内置了 Saga 类,它做的事很简单:注册补偿动作,出异常时按逆序执行。
Saga saga = new Saga(new Saga.Options.Builder()
.setParallelCompensation(false) // 串行补偿,顺序可控
.setContinueWithError(true) // 某个补偿失败也继续后续补偿
.build());
saga.addCompensation(activities::unlockStock, orderId);
activities.lockStock(orderId);
try {
activities.charge(orderId, amount);
} catch (ActivityFailure e) {
saga.compensate();
throw e;
}
三个设计决策要提前想清楚:补偿是并行还是串行(串行更容易排查,并行更快);某个补偿失败是否继续(通常要继续,避免一个卡点导致整体无法回滚);补偿失败后怎么办(记录待人工处理的补偿任务)。
补偿动作本身必须幂等,因为补偿也可能被重复执行。补偿失败时不要吞掉异常,要记录到业务可查询的地方,让运维能看到「这单退款没成功」。这部分与 Saga 与分布式事务补偿 的讨论完全一致。
11. 子工作流与并行执行
复杂流程用子工作流(Child Workflow)拆分,好处是历史独立、可以单独查询、可以复用。与 Activity 的区别是:子工作流是完整的 Workflow,有自己的事件历史与重试语义,且可以返回结构化结果。
// 并行执行三个独立步骤,全部完成后继续
Promise<Quote> flight = Async.function(activities::bookFlight, req);
Promise<Quote> hotel = Async.function(activities::bookHotel, req);
Promise<Quote> car = Async.function(activities::bookCar, req);
activities.confirm(flight.get(), hotel.get(), car.get());
用 Async.function 而不是 CompletableFuture 或线程池,因为前者是确定性的(引擎能重放并行分支),后者不是。这是 Temporal 里最常见的确定性违反。
并行分支的失败处理要显式设计:任何一个 Promise 抛异常,其他 Promise 的结果仍然有效,需要手工补偿已完成的部分。Temporal 不会自动回滚并行分支。
12. 版本兼容与 GetVersion
这是 Temporal 最需要提前规划的部分。Workflow 执行可能跨越数月,期间你部署了新代码,而老实例重放时用的是新代码——如果新代码的控制流变了,重放就会失败(NonDeterminismError)。
正确的做法是用 Workflow.getVersion 做代码分支:
int version = Workflow.getVersion("add-risk-check",
Workflow.DEFAULT_VERSION, 1);
if (version >= 1) {
activities.riskCheck(orderId);
}
activities.charge(orderId, amount);
getVersion 在重放时返回历史记录中的版本号,所以老实例会跳过 riskCheck,新实例会执行它。这个机制要求「变更点是加法而不是改写」:可以新增步骤,但不能删除或调换已有步骤的顺序。
如果不能保持加法(比如要删除一个步骤),替代方案是 Workflow.patched 或者用 ContinueAsNew 开一个新实例。彻底删掉旧版本分支的时机是「所有老实例都已结束」,可以用 temporal workflow list --query "ExecutionStatus='Running'" 检查。
13. ContinueAsNew 与历史截断
事件历史会随执行步数增长。一个每 5 秒轮询一次的 Workflow 跑一天会产生 17280 条事件,历史会膨胀到几 MB,重放变慢,且 Temporal 有 50K 事件或 50 MB 的硬上限。
ContinueAsNew 是解决方案:它结束当前执行,用相同 WorkflowId 启动一个新的执行,把当前状态作为输入传过去。历史从头开始,逻辑继续。
if (Workflow.getInfo().getHistoryLength() > 10_000) {
// 把当前状态传下去,重启历史
continueAsNew(new OrderInput(orderId, currentState));
}
continueAsNew 必须放在 Workflow 代码的最后(它抛出 ContinueAsNewError 终止当前执行)。典型触发条件是「历史长度超过阈值」或「累计处理了 N 条消息」。轮询型 Workflow(长驻、不断处理事件)几乎都必须用这个模式,否则必然撞上历史上限。
注意 ContinueAsNew 与 getVersion 的交互:新执行的第一次重放会从新代码开始,所以版本分支的写法要能兼容「历史为空」的情况。
14. 本地活动与优化手段
本地活动(Local Activity)是不经过服务端调度的 Activity,直接在当前 Worker 进程里执行,结果写进 Workflow 事件历史。它的优势是延迟低(省掉一次往返),劣势是失败后靠 Workflow 重试(而不是 Activity 的重试策略),且历史里会记录完整结果。
适用场景是「执行极快、结果很小、失败率极低」的操作,比如参数校验、格式转换、读取本地缓存。不适用场景是任何有外部副作用的操作——因为本地活动失败后整个 Workflow Task 会重试,可能造成重复执行。
LocalActivityOptions options = LocalActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.build();
其他优化手段:Workflow.newActivityStub 复用(不要每步新建)、Activity 结果尽量小(大结果会写进历史)、把多个小 Activity 合并成一个(减少事件数)。判断标准是「事件数与历史大小」,用 temporal workflow describe 能看到历史长度。
15. 数据转换与序列化
Temporal 默认用 JSON 序列化 Workflow 输入输出与 Activity 参数。支持自定义 DataConverter,比如用 Protobuf 或 Avro:
DataConverter converter = new CompositeDataConverter(
new ProtobufJsonPayloadConverter(),
new JacksonJsonPayloadConverter()
);
WorkflowClient client = WorkflowClient.newInstance(service,
WorkflowClientOptions.newBuilder().setDataConverter(converter).build());
序列化有两个硬约束。第一,类型的向后兼容性:如果你在 Workflow 输入类里删了一个字段,老实例重放时反序列化会失败。所以输入类要遵循「只加字段、不改类型、不删字段」的规则,或者在类上加 @JsonIgnoreProperties(ignoreUnknown = true)。第二,大小限制:单个 Payload 默认上限是 2 MB(可配置),超过要改成传引用(把数据放对象存储,传 key)。
还有一个容易忽略的点:Activity 的返回值也会被写进事件历史。一个返回 1000 条记录的 Activity 会让历史迅速膨胀。正确做法是让 Activity 返回汇总结果,明细走数据库。
16. 与消息中间件的集成
Temporal 与 Kafka 的组合非常常见:Kafka 承载领域事件的扇出,Temporal 承载单笔业务的时序编排。两种集成方向:
- Kafka → Temporal:消费 Kafka 消息,用消息里的业务键作为 WorkflowId 启动或发信号。要保证幂等,重复消费不能重复启动。参见 Kafka 入门 里的消费者语义。
- Temporal → Kafka:Activity 里发消息,用 outbox 模式保证「数据库写入与消息发送」的一致性。
// 消费者里:用业务键做幂等启动
try {
client.newUntypedWorkflowStub("OrderWorkflow", options).start(input);
} catch (WorkflowExecutionAlreadyStarted e) {
// 已存在,改为发信号
client.newUntypedWorkflowStub(existingId).signal("onEvent", event);
}
这种模式的架构意义在 事件驱动架构 里有更完整的讨论。核心原则是「消息负责传输,Workflow 负责状态」,不要让 Kafka 承担状态存储,也不要让 Temporal 承担高吞吐的事件扇出。
17. 与外部系统的幂等桥接
Temporal 保证 Workflow 逻辑的「恰好一次」语义(在事件历史层面),但它无法让外部系统也恰好一次。所以与外部系统交互时必须自己搭幂等桥。
三种常见做法:
- 幂等键传递:Activity 调外部 API 时带上
workflowId + activityId + attempt组成的幂等键,外部系统按此去重。 - 状态表去重:Activity 先查本地状态表,已处理过就直接返回上次结果。
- 查询后写入:先查询外部系统是否已有该笔操作(比如查支付流水),有则跳过。
INSERT INTO activity_dedup (idem_key, result, created_at)
VALUES (?, ?, NOW())
ON DUPLICATE KEY UPDATE result = result;
-- 返回影响行数为 0 说明已存在,读回已有 result 即可
第三种最可靠但依赖外部系统提供查询接口。实践中常常三种混用:本地去重表挡住重复投递,幂等键挡住跨系统重复,查询兜底对账。
18. 部署与运维
Temporal 集群由四个服务组成:Frontend(网关与路由)、History(状态与事件存储)、Matching(任务队列)、Worker(系统内部工作流)。持久化层需要两个数据库:主存储(Cassandra / PostgreSQL / MySQL)和可见性存储(Elasticsearch,用于复杂查询)。
# 单机开发环境(temporal CLI 自带)
temporal server start-dev --db-filename temporal.db --ui-port 8233
生产环境的容量规划要点:History 服务是状态写入的关键路径,要按分片数(shards)规划;主存储的写入 IOPS 决定整体吞吐上限;可见性存储的数据量与保留期决定查询能力。Temporal 默认保留 30 天已关闭工作流的历史,可配置。
升级策略上,Temporal 服务端支持滚动升级,但 SDK 与 Server 有版本兼容矩阵,升级前必须查兼容表。客户端 SDK 升级要小心:新版本可能改变确定性行为,导致老实例重放失败。
19. 自建与 Temporal Cloud 的取舍
| 维度 | 自建 | Temporal Cloud |
|---|---|---|
| 运维成本 | 需要专人维护四个服务与两个数据库 | 零运维 |
| 成本模型 | 服务器成本固定 | 按 Action 计费,量大时贵 |
| 数据合规 | 数据在自有 VPC | 需要评估数据出境 |
| 扩容 | 手工规划分片与存储 | 自动 |
| 版本控制 | 自己决定升级时机 | 跟随官方节奏 |
| 功能差异 | 全功能 | 少数企业特性独占 |
判断标准:如果团队没有专职的平台工程人力,或者业务量不大(每天几十万次 Action 以内),Temporal Cloud 通常更划算,因为自建的隐性成本(人力、故障排查、升级)远高于账单差异。反过来,如果有强数据合规要求或已有成熟的 Cassandra 运维能力,自建更合适。
20. 观测与调试
Temporal 的调试入口是事件历史。Web UI 里能看到每个 Workflow 的完整时间线:什么时候调度了哪个 Activity、重试了几次、失败原因是什么、什么时候收到信号。
排查问题的标准路径是:先看 Workflow 当前状态(Running / Failed / TimedOut),再看最后一个未完成的事件,最后看对应 Activity 的失败原因。如果是 NonDeterminismError,说明代码变更破坏了确定性,需要对照历史找出分支差异。
关键指标要采集四类:Workflow 调度成功率(temporal_workflow_started)、Activity 失败率(temporal_activity_execution_failed)、任务队列积压(temporal_task_queue_lag)、Workflow Task 超时率(temporal_workflow_task_timeout)。任务队列积压是最重要的告警指标,它直接反映 Worker 容量是否足够。
跨 Workflow 的链路追踪需要把 TraceId 存进 Workflow 变量并透传给 Activity,这样 Activity 里的 span 才能挂到同一条链路上,细节见 工作流可观测与调试 。
21. 落地路线图
- 第 1 周:本地起
temporal server start-dev,写一个「Activity + 定时器 + 信号」的最小 Workflow,验证崩溃恢复(执行中杀掉 Worker)。 - 第 2 周:加入重试策略与 Saga 补偿,用一个会随机失败的 Activity 验证重试与回滚。
- 第 3 周:设计 WorkflowId 幂等策略与外部系统幂等桥,做一次重复投递压测。
- 第 4 周:规划版本兼容策略(getVersion 的使用规范)与历史截断(ContinueAsNew 阈值),接入观测指标。
试点流程要选「步骤 5 到 10 步、含一次外部调用、含一次等待」的,不要选最复杂的。验证崩溃恢复最简单的方法是在 Workflow 中间加一个 Workflow.sleep(60s),执行到这里杀掉 Worker,重启后观察是否从断点继续。
22. 权衡取舍
| 选择 | 收益 | 代价 |
|---|---|---|
| 代码即流程 | 无需建模语言,工程师上手快 | 业务方无法直接读流程 |
| 确定性重放 | 崩溃恢复自动化 | 约束多,需要专门规范与评审 |
| Activity 至少一次 | 高可用、易实现 | 必须自己保证幂等 |
| 子工作流拆分 | 历史独立、可复用 | 跨实例调试成本高 |
| getVersion 分支 | 支持长期运行实例 | 代码里积累历史分支,需要清理 |
| ContinueAsNew | 历史可控 | 状态必须能序列化传递 |
| 本地活动 | 低延迟 | 失败重试语义弱,不适合有副作用 |
| Temporal Cloud | 零运维 | 按 Action 计费,成本随量增长 |
23. 常见坑清单
- 在 Workflow 里调
System.currentTimeMillis()或new Date(),本地正常但重放时分支漂移。 - 用
Thread.sleep或CompletableFuture做等待与并行,阻塞 Worker 线程且破坏确定性。 - Activity 未实现幂等,重试时重复扣款,误以为引擎提供恰好一次。
- 只配
maximumAttempts不配nonRetryableErrorTypes,业务错误被无意义重试十次。 - 用随机 UUID 做 WorkflowId,失去幂等启动能力,重复请求产生两个实例。
- 直接修改已有 Workflow 代码的控制流(删步骤、换顺序),老实例重放报
NonDeterminismError。 - 长驻轮询 Workflow 不用 ContinueAsNew,撞上 50K 事件上限后执行失败。
- Activity 返回大对象(上千条明细),事件历史膨胀到几十 MB,重放极慢。
- 输入类删字段导致老实例反序列化失败,没有做向后兼容。
- 所有流程共用一个任务队列,慢 Activity 挤占 Worker 并发槽位。
- 忘记配
startToCloseTimeout,默认 10 分钟把长批处理任务判超时。 - Workflow 代码里打日志,重放时日志重复输出几十遍,淹没真正的错误信息。
24. 小结
Temporal 把「跨天流程的状态管理」这个老问题用一个反直觉的方案解掉了:不存状态,存事件;不恢复状态,重放代码。这个方案的收益是开发者可以继续写普通代码,代价是必须遵守确定性约束,并且要为每一次代码变更考虑运行中实例的兼容性。
落地时的优先级建议是:先解决幂等(WorkflowId 策略 + Activity 幂等键),再解决版本兼容(getVersion 规范),最后解决历史膨胀(ContinueAsNew 阈值)。这三件事里任何一件没做好,都会在上线几个月后以「奇怪的重放错误」或「历史爆掉」的形式暴露。
如果流程里有人工审批环节,Temporal 需要用 Signal 加超时自建一套,可以考虑与 BPMN 引擎分层,参考 人工任务与审批流表单 。如果需要按时间周期批量调度,Temporal 不是合适的工具,应该看 Airflow DAG 调度体系 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。