How we built saga rollbacks for Cloudflare Workflows

TL;DR · AI 摘要
Cloudflare Workflows 现在支持 Saga 回滚功能,允许开发者在步骤内定义补偿逻辑,以处理失败时的事务回退。
核心要点
- Cloudflare Workflows 现在支持 Saga 回滚,允许在步骤内定义补偿逻辑。
- 使用 rollback 参数可以避免手动编写复杂的异常处理逻辑。
- 回滚操作必须是幂等的,并且在失败时继续执行后续操作。
结构提纲
按章节快速跳转。
思维导图
用一张图看清主题之间的关系。
查看大纲文本(无障碍 / 无 JS 友好)
- Cloudflare Workflows Saga 回滚
- Saga 模式
- 补偿机制
- 分布式事务
- 实现方式
- 步骤内定义补偿逻辑
- 避免手动编写异常处理
- 回滚操作
- 幂等性
- 失败时继续执行
金句 / Highlights
值得收藏与分享的关键句。
Saga 模式是一种用于处理分布式事务的补偿机制,确保操作失败时可以回滚。
使用 rollback 参数可以避免手动编写复杂的异常处理逻辑。
回滚操作必须是幂等的,并且在失败时继续执行后续操作。
我们如何为 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() 中定义补偿逻辑,作为步骤本身的一部分,从而保持工作流在回滚时的持久性。
// 记录哪些步骤完成,以便知道需要撤销哪些操作
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;
}没有回滚功能
// 每个步骤都自带撤销操作。添加一个步骤,
// 只需在此处添加其回滚操作。不需要扩展 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() 的最后一个参数传递即可。
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 和执行模型非常重要。
- 失败的步骤可能仍然需要回滚。如果一个步骤注册了回滚处理程序,即使它失败了,它仍然可以回滚。
如果用户代码捕获了错误并且工作流继续执行,回滚不会开始。但如果一个步骤错误被捕获,而工作流之后由于其他原因失败,之前注册的处理程序仍然可以运行,这些处理程序按照步骤启动顺序的逆序执行。
为什么?该步骤可能在失败之前已经部分与外部系统交互。例如,支付提供商可能已经捕获了付款,但该步骤可能在返回 chargeId 给工作流之前失败。这就是为什么回滚处理程序接收 output,但必须处理 output === undefined 的情况。
- 回滚只在工作流失败时开始。添加一个回滚处理程序并不意味着每个步骤错误都会触发回滚。如果用户代码捕获了错误并继续执行,工作流将继续执行。回滚仅在工作流本身即将终端失败时开始。
当回滚开始时,工作流会找到符合条件的 step.do() 调用,运行它们的回滚处理程序,然后记录最终的工作流失败。
- 顺序必须是可预测的。对于顺序工作流,回滚顺序感觉很明显:
- 预留库存。
- 收取信用卡。
- 创建运输。
- 如果运输失败,退款信用卡并释放库存。
并行步骤使这一点更加微妙。完成顺序可能与启动顺序不同,因此工作流使用步骤启动顺序的逆序,而不是完成顺序的逆序。
实际规则是:
- 任何已启动或已完成的步骤,如果有回滚处理程序,都是符合条件的。
- 如果步骤注册了回滚处理程序,失败的步骤.do() 也是符合条件的。
- 处理程序按照步骤启动顺序的逆序执行,而不是完成顺序。
一旦我们明确了预期的行为,就需要将这个新模式添加到 Workflows API 中。回滚功能经历了几次迭代,最终我们确定了回滚选项。
为什么不采用流畅或构建器风格的 API?
最初的方法是采用流畅的风格:step.do(...).rollback(...)。这种写法读起来很自然。向前操作和补偿操作紧挨在一起,调用站点看起来就像普通的 JavaScript 链式调用。
问题是,step.do() 已经有重要的含义:它启动一个持久化的步骤并返回一个步骤输出的 Promise。在 Workers 中,类似 Promise 的值尤其有意义,因为 Workers 的 RPC 支持 Promise 管道(promise pipelining),这是从 Cap'n Proto 等系统继承的模式。
Promise 管道允许代码在值完全返回给调用者之前,就对其调用方法。例如:
const session = api.authenticate(apiKey);
const name = await session.whoami();在这里,session 还不是真正的会话对象。它更像是一个即将出现的会话的句柄。当你调用 session.whoami() 时,Workers 可以提前将这个调用发送到远程端,并告诉它:“一旦认证创建了会话,请在该会话上调用 whoami()。”
这样可以节省一次往返。调用者不需要等待 authenticate() 完全完成后再请求 whoami()。
我们曾考虑过流畅的 API:
step.do("charge-card", chargeCard).rollback(refundCharge);对读者来说,这可能看起来像是“在 charge-card 的结果上调用 .rollback()”。但回滚并不是步骤输出的一部分。它是步骤的选项之一,在步骤开始之前注册的,因此 Workflows 知道如何在后续步骤失败时补偿该步骤。
流畅的 API 还使得步骤的时序更难理解。目前,step.do() 在被调用时立即启动步骤,因此开发者可以启动一个步骤,执行其他操作,并稍后等待第一个步骤的结果:
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:
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(),而不是取代它。
作为步骤元数据的回滚
step.do(..., { rollback })最终,我们选择了显式形式,其中回滚是步骤的元数据。
这样,每个回滚都在正向步骤本身中定义。每个处理程序会接收到导致回滚启动的错误、步骤上下文以及输出,该输出要么是正向步骤返回的持久化值(可以是 undefined),要么是在持久化值之前步骤失败时的 undefined。
回滚会发出生命周期事件,因此你可以知道补偿是否已开始、哪个回滚处理程序失败,以及回滚是否成功完成。
关键的是,原始 Workflow 的失败仍然保持独立:回滚是 Workflow 在失败后执行的操作,而不是 Workflow 失败的原因。
正如你可以在 WorkflowStepConfig 中通过 step 配置定义自定义的重试和超时行为一样,你也可以在 rollbackConfig 中添加回滚相关的特定值。
{
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