The Cloudflare Blog

How we built saga rollbacks for Cloudflare Workflows

8.5内容质量
How we built saga rollbacks for Cloudflare Workflows

TL;DR · AI 摘要

Cloudflare Workflows 现在支持 Saga 回滚功能,允许开发者在步骤内定义补偿逻辑,以处理失败时的事务回退。

核心要点

  • Cloudflare Workflows 现在支持 Saga 回滚,允许在步骤内定义补偿逻辑。
  • 使用 rollback 参数可以避免手动编写复杂的异常处理逻辑。
  • 回滚操作必须是幂等的,并且在失败时继续执行后续操作。

结构提纲

按章节快速跳转。

  1. Cloudflare Workflows 允许构建多步骤的持久化应用,但失败时可能导致状态不一致。

  2. ·Saga 模式

    Saga 模式是一种用于处理分布式事务的补偿机制,确保操作失败时可以回滚。

  3. 开发者现在可以在步骤内定义补偿逻辑,避免手动编写复杂的异常处理逻辑。

  4. 文章提供了带有和不带回滚功能的代码示例,展示了如何实现 Saga 回滚。

  5. 回滚操作必须是幂等的,并且在失败时继续执行后续操作。

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • Cloudflare Workflows Saga 回滚
    • Saga 模式
      • 补偿机制
      • 分布式事务
    • 实现方式
      • 步骤内定义补偿逻辑
      • 避免手动编写异常处理
    • 回滚操作
      • 幂等性
      • 失败时继续执行

金句 / Highlights

值得收藏与分享的关键句。

#Cloudflare#Workflows#Saga#回滚#分布式系统
打开原文

我们如何为 Cloudflare Workflows 构建 saga 回滚

2026-06-25

  • Vaishnav Kavitha
  • Mia Malden
  • André Venceslau

9 分钟阅读

Cloudflare Workflows 允许你构建具有内置重试和状态持久性的多步骤应用程序,适用于长时间运行的进程。当 Workflow 执行时,每一步都可以调用外部系统、重试失败操作,并在重启时持久化状态。但如果某一步失败,可能会导致之前已完成步骤的工作处于不一致或部分状态。

今天,我们发布了 Workflows 的 saga 回滚功能,允许你在步骤内部声明回滚逻辑,以防失败。

例如,考虑一个在两家不同银行之间转账的 Workflow:

  • 从 Bank A 的账户中扣款
  • 存入 Bank B 的账户
  • 向两个账户的所有者发送电子邮件确认

如果第二步(向 Bank B 的账户存入资金)失败了会怎样?一旦从 Bank A 扣款成功,交易就会提交,资金已经离开了该系统。作为交易的协调者,你不能简单地在 Bank A 的系统中“撤销”该操作。相反,必须通过一个新的操作将资金重新存入 Bank A 的账户,该操作在语义上与第一个操作相反。

这种操作与其补偿逻辑的配对被称为 saga 模式。

在今天之前,开发者必须自己实现补偿逻辑,以跟踪哪些操作成功、哪些失败,以及在失败时应采取哪些操作,这些都必须在步骤的直接定义之外进行。现在,你可以在每个步骤的 do() 中定义补偿逻辑,作为步骤本身的一部分,从而保持工作流在回滚时的持久性。

code
// 记录哪些步骤完成,以便知道需要撤销哪些操作
let debitA;
let creditB;
try {
  debitA = await step.do("debit-bank-a", () => bankA.debit(from, amount));
  creditB = await step.do("credit-bank-b", () => bankB.credit(to, amount));
  await step.do("notify", () => notifyBoth(from, to, amount));
} catch (error) {
  // 逆序回滚。每个撤销操作本身都是一个持久化的步骤,
  // 必须是幂等的,并且如果其中一个失败,必须继续执行。
  if (creditB) {
    try {
      await step.do("reverse-credit-b", () => bankB.debit(to, amount, creditB.id));
    } catch (e) {
      await alertOnCall("reverse-credit-b failed", e);
    }
  }
  if (debitA) {
    try {
      await step.do("refund-debit-a", () => bankA.credit(from, amount, debitA.id));
    } catch (e) {
      await alertOnCall("refund-debit-a failed", e);
    }
  }
  throw error;
}

没有回滚功能

code
// 每个步骤都自带撤销操作。添加一个步骤,
// 只需在此处添加其回滚操作。不需要扩展 catch
// 块,不需要手动排序,不需要重播逻辑。
await step.do("debit-bank-a", () => bankA.debit(from, amount), {
  rollback: async ({ output }) => bankA.credit(from, amount, output.id),
});
await step.do("credit-bank-b", () => bankB.credit(to, amount), {
  rollback: async ({ output }) => bankB.debit(to, amount, output.id),
});
await step.do("notify", () => notifyBoth(from, to, amount));

带有回滚功能

尝试使用

要使用回滚功能,只需将一个包含回滚函数的选项对象作为 step.do() 的最后一个参数传递即可。

code
const debit = await step.do(
  "debit-account-a",
  async () => {
    return await bankA.debit({
      accountId: fromAccountId,
      amount,
      idempotencyKey: `${transferId}:debit-account-a`,
    });
  },
  {
    rollback: async () => {
      await bankA.credit({
        accountId: fromAccountId,
        amount,
        idempotencyKey: `${transferId}:rollback-debit-account-a`,
      });
    },
  }
);

// 唯一性键确保了正向操作和回滚操作都可以安全重试,而不会重复转账

const credit = await step.do(
  "credit-account-b",
  async () => {
    return await bankB.credit({
      accountId: toAccountId,
      amount,
      idempotencyKey: `${transferId}:credit-account-b`,
    });
  },
  {
    rollback: async ({ output }) => {
      if (output === undefined) {
        return;
      }

      await bankB.debit({
        accountId: toAccountId,
        amount,
        idempotencyKey: `${transferId}:rollback-credit-account-b`,
      });
    },
  }
);

// 如果在此处失败,我们可能需要撤销所有之前的付款。用户不应该需要将代码包裹在复杂的 try-catch 逻辑中,只是为了撤销两次小的付款(见下文)

await step.do("send-confirmation", async () => {
  await sendTransferConfirmation({ ... });
});

回滚函数应该和常规的工作流步骤一样是幂等的。如果你退款,使用支付提供商的唯一性键。如果你释放库存,请确保释放操作可以多次调用而不会出现问题。

如果任何步骤失败,回滚处理程序将按照与步骤启动顺序相反的顺序执行。听起来很简单:当某事失败时,运行撤销步骤。实际上,有一些细节使得 API 和执行模型非常重要。

  1. 失败的步骤可能仍然需要回滚。如果一个步骤注册了回滚处理程序,即使它失败了,它仍然可以回滚。

如果用户代码捕获了错误并且工作流继续执行,回滚不会开始。但如果一个步骤错误被捕获,而工作流之后由于其他原因失败,之前注册的处理程序仍然可以运行,这些处理程序按照步骤启动顺序的逆序执行。

为什么?该步骤可能在失败之前已经部分与外部系统交互。例如,支付提供商可能已经捕获了付款,但该步骤可能在返回 chargeId 给工作流之前失败。这就是为什么回滚处理程序接收 output,但必须处理 output === undefined 的情况。

  1. 回滚只在工作流失败时开始。添加一个回滚处理程序并不意味着每个步骤错误都会触发回滚。如果用户代码捕获了错误并继续执行,工作流将继续执行。回滚仅在工作流本身即将终端失败时开始。

当回滚开始时,工作流会找到符合条件的 step.do() 调用,运行它们的回滚处理程序,然后记录最终的工作流失败。

  1. 顺序必须是可预测的。对于顺序工作流,回滚顺序感觉很明显:
  • 预留库存。
  • 收取信用卡。
  • 创建运输。
  • 如果运输失败,退款信用卡并释放库存。

并行步骤使这一点更加微妙。完成顺序可能与启动顺序不同,因此工作流使用步骤启动顺序的逆序,而不是完成顺序的逆序。

实际规则是:

  • 任何已启动或已完成的步骤,如果有回滚处理程序,都是符合条件的。
  • 如果步骤注册了回滚处理程序,失败的步骤.do() 也是符合条件的。
  • 处理程序按照步骤启动顺序的逆序执行,而不是完成顺序。

一旦我们明确了预期的行为,就需要将这个新模式添加到 Workflows API 中。回滚功能经历了几次迭代,最终我们确定了回滚选项。

为什么不采用流畅或构建器风格的 API?

最初的方法是采用流畅的风格:step.do(...).rollback(...)。这种写法读起来很自然。向前操作和补偿操作紧挨在一起,调用站点看起来就像普通的 JavaScript 链式调用。

问题是,step.do() 已经有重要的含义:它启动一个持久化的步骤并返回一个步骤输出的 Promise。在 Workers 中,类似 Promise 的值尤其有意义,因为 Workers 的 RPC 支持 Promise 管道(promise pipelining),这是从 Cap'n Proto 等系统继承的模式。

Promise 管道允许代码在值完全返回给调用者之前,就对其调用方法。例如:

javascript
const session = api.authenticate(apiKey);
const name = await session.whoami();

在这里,session 还不是真正的会话对象。它更像是一个即将出现的会话的句柄。当你调用 session.whoami() 时,Workers 可以提前将这个调用发送到远程端,并告诉它:“一旦认证创建了会话,请在该会话上调用 whoami()。”

这样可以节省一次往返。调用者不需要等待 authenticate() 完全完成后再请求 whoami()

我们曾考虑过流畅的 API:

javascript
step.do("charge-card", chargeCard).rollback(refundCharge);

对读者来说,这可能看起来像是“在 charge-card 的结果上调用 .rollback()”。但回滚并不是步骤输出的一部分。它是步骤的选项之一,在步骤开始之前注册的,因此 Workflows 知道如何在后续步骤失败时补偿该步骤。

流畅的 API 还使得步骤的时序更难理解。目前,step.do() 在被调用时立即启动步骤,因此开发者可以启动一个步骤,执行其他操作,并稍后等待第一个步骤的结果:

javascript
const first = step.do("first", () => serviceA.call());

await step.do("second", () => serviceB.call());

await first;

在目前的执行模型中,first 会立即启动,而不是在 second 之后。流畅的 API 会增加复杂性。Workflows 需要等待并查看是否附加了 .rollback(),才能知道完整的步骤定义。这可能会延迟步骤发送到引擎的时间。

在前面的例子中,first 可能会在 await first 时启动,而不是在 step.do("first", ...) 时,这发生在 second 已经完成之后。

这使得并发的 Workflows 更难理解:步骤的时序将取决于返回的 Promise 被消费的时间,而不仅仅是 step.do() 被调用的位置。

我们还考虑过构建器风格的 API:

javascript
const charge = await step
	.saga("charge")
	.do(() => chargeCard())
	.rollback(() => refundCharge())
	.run();

构建器风格的 API 避免了 Promise 的歧义。它还为我们提供了一个明显的位置来添加未来步骤级别的选项,并清楚地表明向前操作和回滚操作属于同一个 Saga 步骤。

但它增加了仪式感。每个步骤都需要一个最终的 .run(),忘记 .run() 会很容易,而且在没有工具支持的情况下很难发现。简单的单步骤情况开始看起来像是配置链。它还引入了一个新的 step.saga() 构建器,与现有的 step.<action> 模式不一致。最重要的是,它让 step.do() 看起来像是一个较旧的 API,而不是 Workflows 的主要原语。回滚的目标是扩展 step.do(),而不是取代它。

作为步骤元数据的回滚

javascript
step.do(..., { rollback })

最终,我们选择了显式形式,其中回滚是步骤的元数据。

这样,每个回滚都在正向步骤本身中定义。每个处理程序会接收到导致回滚启动的错误、步骤上下文以及输出,该输出要么是正向步骤返回的持久化值(可以是 undefined),要么是在持久化值之前步骤失败时的 undefined。

回滚会发出生命周期事件,因此你可以知道补偿是否已开始、哪个回滚处理程序失败,以及回滚是否成功完成。

关键的是,原始 Workflow 的失败仍然保持独立:回滚是 Workflow 在失败后执行的操作,而不是 Workflow 失败的原因。

正如你可以在 WorkflowStepConfig 中通过 step 配置定义自定义的重试和超时行为一样,你也可以在 rollbackConfig 中添加回滚相关的特定值。

code
{
  rollback: async ({ output }) => {
    await bankA.credit({ accountId: fromAccountId, amount, transferId: `${transferId}-reversal` });
  },
  rollbackConfig: {
    retries: { limit: 10, delay: '30 seconds', backoff: 'exponential' },
    timeout: '2 minutes',
  },
}

这与我们想要的生命周期事件心理模型一致。step.do() 已经描述了一个 Workflows 记录、重试并在日志中稍后显示的持久工作单元。回滚是该相同工作单元的另一个生命周期行为。它应该与步骤定义一起旅行,而不是存在于单独的包装器或构建器中。

  • 步骤仍然在 step.do() 正常启动时开始。
  • 返回的 Promise 仍然代表步骤的输出。
  • 并发的 Workflow 代码保持相同的执行模型。
  • 回滚的重试和超时选项与回滚处理程序相邻。
  • 现有的 step.do() 调用仍然像今天一样完全正常工作。

这种结构比流畅的 API 稍微更显式一些,但这种显式性是有用的。操作及其补偿仍然在一个地方,API 不会引入新的步骤构建器或新的 Promise 类型。已经理解 step.do() 的开发人员只需要学习一个额外的选项对象。

这看起来不那么神奇,但更容易采用,也更容易理解。

内部工作原理

回滚看起来像是一个小小的 API 增加,但它改变了 Workflows 需要为每个步骤记录的内容。

一个普通的 step.do() 已经有一个持久的记录。Workflows 记录该步骤是否开始、是否完成、返回了什么内容,以及如果 Workflow 后续恢复,是否应该跳过而不是重复该步骤。

回滚则在该记录中增加了一项内容:该步骤是否注册了补偿逻辑。

这意味着,如果 Workflow 失败,Workflows 需要整合两部分信息。

第一部分是持久的步骤历史。Workflow 引擎存储数据以了解哪些步骤运行了、哪些步骤完成、保存了什么输出,以及是否注册了回滚。

第二部分是回滚处理程序本身,即用于补偿该步骤的函数。Workflows 不会将该函数的文本内容保存为数据。相反,它在 Workflow 运行期间保留对该处理程序的可调用引用。

在 Workers RPC 中,这种可调用的引用被称为 stub(存根)。stub 允许系统的一部分调用在其他地方运行的代码。stub 还具有生命周期,当调用或执行上下文结束时,它们可以被丢弃。如果你需要在该点之后继续使用 stub,Workers RPC 提供了一个 dup() 方法,它可以创建另一个指向相同目标的句柄。

对于回滚来说,这种模型非常有用。持久化的步骤历史记录了需要补偿的内容。回滚 stub 为 Workflows 提供了一种调用补偿代码的方式。由于回滚处理程序可能需要在注册它们的 step.do() 调用之后继续存在,Workflows 会为回滚阶段保留一个可调用的引用指向处理程序。

在常见情况下,当 Workflow 在相同的引擎生命周期内进入回滚时,Workflows 已经拥有它所需的回滚 stub。它可以使用持久化的步骤历史找到符合条件的步骤,然后调用在正向执行期间注册的回滚 stub。

当 Workflows 需要在重启后恢复时,情况会变得更加微妙。

如果引擎在需要回滚时被驱逐、崩溃或重启,Workflows 仍然拥有持久化的步骤历史,但它可能不再拥有内存中的回滚 stub。为了恢复,Workflows 使用重放(replay):一种恢复模式,它可以在不重新执行已完成的正向步骤体的情况下重新运行 Workflow 代码。

当重放到达一个已完成的 step.do() 时,Workflows 会读取持久化的结果,而不是再次运行步骤体。对于回滚恢复,Workflows 只需要为那些附加了回滚并符合回滚条件的步骤重新构建处理程序。当这些 step.do() 调用被遇到时,它们的回滚选项可以再次注册可调用的 stub。

这使得 Workflows 能够在不复制原始外部副作用的情况下恢复所需的回滚处理程序。

有了这些组件,无论处理程序是否仍然在内存中可用,或者是否需要在恢复期间重新构建,回滚都可以正常工作。

当工作流即将失败时,Workflows 不会要求你的应用程序重建发生的事情。它已经拥有步骤历史。它可以查看持久化的记录并回答重要的问题:

  • 哪些步骤已经启动了?
  • 哪些步骤已经完成?
  • 哪些失败的步骤可能仍然需要清理?
  • 哪些步骤注册了回滚处理程序?
  • 每个回滚处理程序应该接收什么输出?
  • 补偿应该以什么顺序运行?

然后,Workflows 会使用回滚上下文调用每个回滚 stub:原始错误、步骤上下文,以及如果有的话,步骤输出。

顺序的细节很重要。在正常的 JavaScript 中,尤其是使用 Promise.all() 时,完成顺序并不总是与启动顺序一致。如果步骤 A 先启动,步骤 B 后启动,步骤 B 可能会先完成。对于回滚,Workflows 使用持久化的启动顺序作为稳定的事实来源,然后按相反顺序进行回滚。

回滚处理程序也通过 Workflows 的正常步骤机制运行。这意味着补偿将具有你期望的 Workflows 的相同操作属性:重试、超时、生命周期事件、日志,以及最终的记录结果。如果回滚处理程序在配置的重试次数后仍然持续失败,Workflows 会将回滚结果记录为失败,停止运行剩余的回滚处理程序,并最终将 Workflow 实例置于 Errored(错误)状态。

这是 saga 回滚与 catch 块之间的主要区别。catch 块只知道在 JavaScript 执行的特定时刻内存中仍然存在的内容。而工作流回滚使用持久化的步骤历史来决定已经发生了什么,通常会调用它已经拥有的存根(stub),并在恢复过程中安全地重建缺失的存根。

这也是为什么 API 将回滚直接放在 step.do() 本身的原因。回滚不是一个独立的全局错误处理程序,而是附加到工作流已经理解的持久工作单元的元数据。

下一步是什么

我们第一版的回滚功能包括:

  • 为 step.do() 提供每个步骤的显式回滚处理程序
  • 顺序执行回滚
  • 补偿操作的重试和超时配置

接下来,我们希望探索:

  • 对 waitForEvent 的回滚支持
  • 并行回滚执行的支持
  • 对 Python 工作流的回滚支持

当一个多步骤的应用程序在执行中途失败时,最难的部分通常不是知道它已经失败了,而是知道已经发生了什么,以及接下来需要发生什么。

Saga 回滚允许你将答案直接放在每个步骤旁边。如果你正在使用工作流构建多步骤的应用程序,请尝试使用 saga 回滚,并告诉我们你接下来想要的补偿模式。从工作流文档开始,然后在 Cloudflare 社区分享你的反馈。

[if astro]>server-island-start<![endif]

Workflows

Cloudflare Workers

Developers