Temporal 与持久化执行

本文深入讲解 Temporal 的持久化执行模型,回答代码即流程如何做到崩溃不丢状态、确定性约束到底限制了什么、版本升级如何不破坏运行中实例。覆盖事件溯源重放、Workflow 与 Activity 职责划分、任务队列与 Worker、信号与定时器、重试策略、Saga 补偿、ContinueAsNew 历史截断与版本兼容,并给出 Java 与 Python 的真实代码片段。

引言

传统写法里,一个「下单后等支付、支付后发货、超时未支付则取消」的流程会被拆成若干个回调接口加一张状态表,再加一个定时任务扫超时。这套东西能跑,但每加一个步骤,状态分支就翻一倍,最终没人能说清「订单处于状态 X 时收到事件 Y 应该怎么办」。

Temporal 提出的解法叫持久化执行(Durable Execution):你写的就是一段普通的、同步的、带 if-else 和 try-catch 的代码,引擎把这段代码的每一次副作用记录成事件,进程崩溃后从事件历史重放,把执行恢复到崩溃前的状态。于是「等待 30 天」在代码里就是 Workflow.sleep(30 天),「失败重试」就是给 Activity 配一个 RetryPolicy,不需要状态表,不需要回调接口。

代价是引入了一个反直觉的约束:Workflow 代码必须是确定性的。同一段代码在重放时必须走完全相同的分支,否则状态会漂移。这个约束是 Temporal 所有「坑」的根源,也是理解它的关键。

本文先讲清楚重放模型到底是怎么工作的,再逐条拆解确定性约束,然后覆盖 Workflow 与 Activity 的职责划分、任务队列、信号、定时器、重试、Saga、版本兼容与历史截断,最后给出部署运维与落地路线。想先看整体选型框架的读者,可以从 工作流引擎全景与选型 开始。

目录

  1. 持久化执行要解决什么问题
  2. 事件溯源与重放模型
  3. 确定性约束清单
  4. Workflow 与 Activity 的职责划分
  5. 一个完整的订单工作流
  6. Worker 与任务队列
  7. 启动、查询与信号
  8. 定时器与长睡眠
  9. 重试策略与超时
  10. Saga 补偿模式
  11. 子工作流与并行执行
  12. 版本兼容与 GetVersion
  13. ContinueAsNew 与历史截断
  14. 本地活动与优化手段
  15. 数据转换与序列化
  16. 与消息中间件的集成
  17. 与外部系统的幂等桥接
  18. 部署与运维
  19. 自建与 Temporal Cloud 的取舍
  20. 观测与调试
  21. 落地路线图
  22. 权衡取舍
  23. 常见坑清单
  24. 小结

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. 常见坑清单

  1. 在 Workflow 里调 System.currentTimeMillis() 或 new Date(),本地正常但重放时分支漂移。
  2. 用 Thread.sleep 或 CompletableFuture 做等待与并行,阻塞 Worker 线程且破坏确定性。
  3. Activity 未实现幂等,重试时重复扣款,误以为引擎提供恰好一次。
  4. 只配 maximumAttempts 不配 nonRetryableErrorTypes,业务错误被无意义重试十次。
  5. 用随机 UUID 做 WorkflowId,失去幂等启动能力,重复请求产生两个实例。
  6. 直接修改已有 Workflow 代码的控制流(删步骤、换顺序),老实例重放报 NonDeterminismError。
  7. 长驻轮询 Workflow 不用 ContinueAsNew,撞上 50K 事件上限后执行失败。
  8. Activity 返回大对象(上千条明细),事件历史膨胀到几十 MB,重放极慢。
  9. 输入类删字段导致老实例反序列化失败,没有做向后兼容。
  10. 所有流程共用一个任务队列,慢 Activity 挤占 Worker 并发槽位。
  11. 忘记配 startToCloseTimeout,默认 10 分钟把长批处理任务判超时。
  12. Workflow 代码里打日志,重放时日志重复输出几十遍,淹没真正的错误信息。

24. 小结

Temporal 把「跨天流程的状态管理」这个老问题用一个反直觉的方案解掉了:不存状态,存事件;不恢复状态,重放代码。这个方案的收益是开发者可以继续写普通代码,代价是必须遵守确定性约束,并且要为每一次代码变更考虑运行中实例的兼容性。

落地时的优先级建议是:先解决幂等(WorkflowId 策略 + Activity 幂等键),再解决版本兼容(getVersion 规范),最后解决历史膨胀(ContinueAsNew 阈值)。这三件事里任何一件没做好,都会在上线几个月后以「奇怪的重放错误」或「历史爆掉」的形式暴露。

如果流程里有人工审批环节,Temporal 需要用 Signal 加超时自建一套,可以考虑与 BPMN 引擎分层,参考 人工任务与审批流表单 。如果需要按时间周期批量调度,Temporal 不是合适的工具,应该看 Airflow DAG 调度体系 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「工作流引擎」更多文章

  1. 工作流成本优化
  2. 执行器与资源隔离
  3. 调度、回填与补数