返回
RSS InfoQ (AI/数据工程) AI 逐段翻译 发布 2026-09-14 19:00

无需外部编排器:在 Postgres 上实现持久工作流

DataHot 速览

文章介绍如何用 Postgres 充当持久状态存储与协调层,从而避免引入 Temporal 或 AWS Step Functions 等外部编排器。它用 SELECT ... FOR UPDATE SKIP LOCKED 构建并发工作队列,用主键检查点强制幂等,并通过租约与清理器实现崩溃恢复。工作流休眠和人工审批也可持久化为数据库状态,在重启后继续执行。作者来自 Kestrel Workflows,示例场景包括事故响应、云资源供给、CI/CD 和开发者请求。

为什么值得关注:为需要在数据平台或自动化系统中实现可靠工作流的从业者提供数据库原生方案,减少外部依赖,并覆盖幂等、并发与崩溃恢复等关键机制。

本文目录 11 节
  1. 关键要点
  2. 默认做法:使用外部编排器
  3. 将 Postgres 作为编排器
  4. 一个让一切变得安全的子句
  5. 让数据库强制执行幂等性
  6. 崩溃恢复:租约与清扫器
  7. 持久休眠与人工审批
  8. 为什么这能扩展
  9. 运维方面的考虑
  10. 免费的观测能力
  11. 可靠性和安全性坍缩为一个依赖项

译文

AI 逐段翻译

关键要点

  • 持久执行并不需要像 Temporal 或 AWS Step Functions 这样的专用编排器;你已经在运维的关系型数据库可以扮演这个角色,并将一个有状态的外部系统从你的关键路径中移除。
  • SELECT ... FOR UPDATE SKIP LOCKED 将一张普通的 Postgres 表变成一个并发工作队列,每一行都恰好被一个工作进程认领,无需消息代理、领导者选举或外部锁服务。
  • 将步骤检查点设为主键约束,可以让数据库强制执行幂等性,这样在崩溃后重新运行的步骤会读取其先前的结果,而不是重复产生副作用。
  • 崩溃恢复可以通过一个简单的租约与清扫器模式来实现,工作进程为其拥有的行发送心跳,而一个周期性查询会重新入队任何租约已过期的执行。
  • 将工作流状态存储在你的主数据库中,使可观测性变成一条简单的 SQL 查询,并将可靠性和安全性收敛到单一依赖上。

当我们开始构建 Kestrel Workflows(用于自动化事件响应、云资源调配、CI/CD 和开发者请求的确定性工作流)时,我们遇到了每个自动化系统最终都会遇到的问题。我们的工作流必须持久,能够经受各种类型的故障。

一个工作流可以被定义为在 Kubernetes 工作负载失败时触发,运行 AI 根因分析,生成修复方案,等待人工通过 Slack 批准修复,然后打开一个 GitOps 拉取请求。如果我们的任何系统在这样一个工作流执行过程中被重新调度(在 Kubernetes 中,这极有可能发生),我们不能就这样丢失工作流执行。我们需要持久执行。

持久执行背后的理念非常简单而强大。当程序运行时,它会将进度检查点保存到像数据库这样的外部系统。如果进程终止,另一个进程会从外部系统重新加载最后一个检查点,并从最后完成的步骤继续。这就像电子游戏中的自动保存,但用于后端代码。

默认做法:使用外部编排器

实现持久执行的通常建议是使用外部编排器,比如 TemporalAWS Step Functions,它们充当协调工作流中各个步骤的中央编排器。在这种模型中,客户端提交一个工作流,编排器将其持久化到自己的数据存储中,然后将其分派给一个工作进程。

对于每个完成的步骤,工作流将其结果发送回编排器,编排器检查点记录进度,然后分派下一步。如果工作进程死亡,编排器会将其步骤重新分配给一个健康的工作进程。

图 1. 外部编排,其中专用编排器协调分派、队列、检查点和工作者。(来源:作者创建。)

我们几乎走了这条路,但在审视了它会给我们带来的代价后重新考虑。

使用外部编排器意味着要再部署、保护、监控和升级另一个有状态系统。它位于每个工作流的关键路径上,因此成为单点故障。它需要自己的访问控制和审计处理,因为工作流步骤负载——对我们来说包括基础设施拓扑、源代码和日志等内容——现在会传输到外部系统。它带来了自己的数据模型,通常是针对工作流状态优化的键值存储,而不是任意关系查询,使得临时分析不如直接查询 Postgres 那样直接。

但事情是这样的:持久执行本质上关乎数据库,因为它就是在将状态检查点保存到一个持久的外部系统。我们已经将 Postgres 作为我们的记录系统运行,并投入了大量精力确保它能够大规模地运行,所以我们问自己:为什么要使用外部编排器呢?为什么不让 Postgres 成为编排器?

图 2. 由 Postgres 支持的持久执行,其中无状态应用服务器直接通过 Postgres 协调工作流状态。(来源:作者创建。)

将 Postgres 作为编排器

在 Kestrel 中,没有编排器进程。每个应用服务器都运行一个嵌入式持久工作流库,并直接与 Postgres 通信。每一个触发器,从 PagerDuty 告警到 GitHub webhook,都会向 workflow_executions 表插入一行:

CREATE TABLE workflow_executions (
	id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
 	workflow_name TEXT NOT NULL,
 	status TEXT NOT NULL DEFAULT 'enqueued',
 	input JSONB NOT NULL,
	owner_id TEXT, -- which worker holds the lease
	lease_expires TIMESTAMPTZ, -- when the lease lapses
	created_at TIMESTAMPTZ NOT NULL DEFAULT now(), 
	updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);

服务器轮询此表以认领工作。对于每个完成的步骤,服务器会将输出检查点保存到 Postgres 本身。

一个让一切变得安全的子句

所有协调都通过 Postgres 数据库进行。使其安全的诀窍是一个子句:

SELECT id, input
FROM workflow_executions
WHERE status = 'enqueued'
ORDER BY created_at
LIMIT 1
FOR UPDATE SKIP LOCKED;

FOR UPDATE SKIP LOCKED 是大多数由 Postgres 支持的队列背后的无名英雄。它锁定工作进程抓取的行,并告诉其他所有工作进程跳过这些行而不是阻塞,从而为系统提供恰好一次的处理保证。两个服务器可以同时轮询同一个表,而绝不会将同一个工作流交给两个执行器。

尽管本文重点关注 Postgres,但这种队列模式并非 Postgres 所独有:MySQL 8.0+MariaDB 10.6+Oracle 也支持 SKIP LOCKED,而 Db2SQL Server 提供了跳过锁定行的等效机制。

这个认领发生在一个短事务内,该事务同时翻转状态并写入租约,因此工作进程在开始任何工作之前就拥有了该行:

BEGIN;
 
WITH claimed AS (
	SELECT id
	FROM workflow_executions
	WHERE status = 'enqueued'
	ORDER BY created_at
	LIMIT 1
	FOR UPDATE SKIP LOCKED
)
UPDATE workflow_executions e
SET status = 'running',
     	owner_id = $1, -- the worker's id
     	lease_expires = now() + INTERVAL '30 seconds',
     	updated_at = now()
FROM claimed
WHERE e.id = claimed.id
RETURNING e.id, e.input;
COMMIT;

这个事务就是整个调度器。你不需要代理、领导者选举或带 TTL 的 Redis 锁。在生产实现中有几件事应该被明确指出:

  • 认领事务应该很小。它应该只锁定行、翻转状态并提交。如果你在步骤运行时保持 FOR UPDATE 事务打开,你将在可能长时间运行的外部工作期间一直占用行锁。
  • LIMIT 可以大于 1,以认领一批执行并分摊往返开销,但批量越大,工作进程在完成它们之前死亡时的影响范围就越大。
  • SKIP LOCKED 在争用情况下不保持 FIFO;对于我们的工作流来说这没问题,而且实际上是我们想要的,但如果你需要严格排序,这种模式就不适用。

让数据库强制执行幂等性

步骤检查点利用第二个 Postgres 原语:完整性约束。每个步骤将其输出写入一个 operation_outputs 表,以 (execution_id, step_id) 作为主键:

CREATE TABLE operation_outputs (
	execution_id BIGINT NOT NULL REFERENCES workflow_executions(id),
	step_id TEXT NOT NULL,
	output JSONB NOT NULL,
	created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
	PRIMARY KEY (execution_id, step_id)
);

当一个步骤完成时,worker 通过一个从不覆盖的 upsert 记录结果:

INSERT INTO operation_outputs (execution_id, step_id, output)
VALUES ($1, $2, $3)
ON CONFLICT (execution_id, step_id) DO NOTHING
RETURNING output;

如果恢复中的 worker 重新运行一个已经提交的步骤,唯一约束会阻止重复的检查点被插入。worker 改为读取现有的检查点并返回先前的。结果。是数据库,而不是应用程序代码,强制实现了幂等性。

在执行一个步骤之前,运行在 worker 上的 Orchestrator SDK 只是检查 operation_outputs 中是否存在 (execution_id, step_id)。如果存在一行,它就跳过执行并将存储的输出向下游返回。在这里,我们再次将持久执行的一个不变量表达为约束,并让 Postgres 原语来维护它,而不是从零开始编写分布式系统的簿记逻辑。

崩溃恢复:租约与清扫器

有意思的故障是 worker 认领了一个执行、将状态设置为 running,然后死掉的情况,因为 OOM 终止、节点排空、Pod 驱逐,或十几种其他 Kubernetes 故障原因。这些行会卡在 running 状态而没有活跃的拥有者;这正是租约列发挥作用的地方。

拥有一个 running 执行的 worker 会按间隔对其发送心跳,将租约向前推进:

UPDATE workflow_executions
SET lease_expires = now() + INTERVAL '30 seconds',
	updated_at = now()
WHERE id = $1 and owner_id = $2;

用 Postgres 实现恢复与强制并发和幂等性一样简单。一个清扫器会重置那些 worker 未续约租约的执行,而一个健康的 worker 会重新加载检查点并从最近完成的步骤继续:

UPDATE workflow_executions
SET status = 'enqueued', owner_id = NULL
WHERE status = 'running'
AND lease_expires < now();

在生产使用中实现这一点时,有两个重要细节。租约时长是一种权衡。如果太短,一个缓慢但存活的 worker 的工作可能被窃取,导致两个 worker 运行同一个执行(得益于幂等的检查点,这是安全的,但仍然是浪费)。如果太长,那么崩溃后的恢复就会被延迟。

如果没有幂等检查点——这在 Postgres 中是免费的——我们就无法使用激进的租约,因为一次误报的清扫会导致带有副作用的双重执行,而不仅仅是重复的工作。

持久休眠与人工审批

即便是等待人类从 Slack 审批工作流这种可能耗时任意长的时间,也是持久的。它只是向一个 workflow_waits 表插入一行,带有一个 wake_at 时间戳,而不是让 goroutine 祈祷它所在的 Pod 能活得足够久以至不会错过审批决定:

CREATE TABLE workflow_waits (
	execution_id BIGINT NOT NULL REFERENCES workflow_executions(id),
	step_id TEXT NOT NULL,
	wake_at TIMESTAMPTZ NOT NULL,
	PRIMARY KEY (execution_id, step_id)
);

当一个工作流执行遇到休眠或审批闸门时,它将自身状态设置为 waiting,并插入一行 workflow_waits 记录,带有一个 wake_at 时间戳,同时一个周期性清扫器会重新入队任何到期的事项:

UPDATE workflow_executions e
SET status = 'enqueued'
FROM workflow_waits w
WHERE w.execution_id = e.id
	AND e.status = 'waiting'
	AND w.wake_at <= now();

这样,等待和审批工作流步骤就能在任意次数的重启或重新调度中存活,因为它们的状态被持久地存储在 Postgres 中,而不是内存中的线程。无论等待或审批窗口有多长,无论是两小时还是两周,成本都是一样的:只有一行。

为什么这能扩展

当你把编排推进到 Postgres 时,可扩展性和可用性是难题,这是件好事,因为它们已被充分理解并有经过验证的解决方案。在 Kestrel 的案例中,这些是我们之前已经解决过的问题。我们通过添加无状态 worker 来扩展吞吐量;瓶颈变成了 Postgres 能多快地清空队列。我们已经看到,单个实例每秒可以处理数万个工作流,而无需添加只读副本或使用 Citus 之类的工具。

可用性就变成了你的 Postgres 高可用设置已经提供的东西,例如流复制和托管的多可用区故障转移,因为 worker 是可互换的,并且可以恢复彼此的状态。

运维方面的考虑

连接压力是需要注意的,因为数百个轮询 worker,每个都持有少量连接,会很快耗尽 Postgres 的连接限制。解决方案之一是在 Postgres 前面使用 PgBouncer 的事务池化,并让空闲 worker 以带抖动的轮询间隔退避。但如果你需要跨数千个 worker 的扇出和亚毫秒级的调度延迟,像 Temporal 这样的外部持久执行引擎可能是更好的架构选择。因为我们的工作负载是 I/O 密集型的自动化流水线,运行时间从几秒到几小时不等,只需要适度的并发,并且大量依赖外部调用,使用 Postgres 不是一种妥协,它实际上是更合适的选择。

另一件需要留意的事情是死元组在高变更率下成为瓶颈。一个繁忙的队列表可能成为 UPDATEDELETE 的热点,从而产生膨胀。幸运的是,在这些表上调优 autovacuum、用部分索引保持热行集很小,以及将已完成的执行分区离开活跃队列,将使工作集保持很小并完全避免这个问题。

免费的观测能力

观测能力现在是免费的,这是我在使用 Postgres 时最喜欢的一点。每个工作流和每个步骤都是表中的一行,所以监控就是 SQL。例如,如果你想要本月所有出错的运行,你只需要执行:

SELECT * FROM workflow_executions
WHERE created_at > NOW() - INTERVAL '1 month'
AND status = 'error';

这看起来极其简单,但它之所以成为可能,只是因为 Postgres 是关系型的。你可以做更强大的事情,比如将执行与步骤和审批连接起来,回答诸如“哪些修复被拒绝了,在哪个步骤被拒绝?”之类的问题。

我们还在 status 和 creation time 等列上维护了一些有针对性的二级索引,以保持这些查询很快,同时不增加太多写入开销。这些索引让我们能够构建一个实时工作流观测仪表盘,而无需额外基础设施。使用外部编排器时,这类关系查询需要将工作流状态从键值存储导出到单独的分析系统。

可靠性和安全性坍缩为一个依赖项

使用 Postgres 实现持久工作流的另一个好处是,可靠性和安全性坍缩为单一依赖。唯一能让 Kestrel 工作流宕机的事情就是 Postgres 故障,但作为我们的主要记录系统,如果 Postgres 故障,其他一切也都会故障。在工作流量较高时,将工作流状态与应用状态隔离到单独的 Postgres 部署中以实现资源隔离是合理的,同时保持相同的运维模型。由于没有任何工作流数据离开数据库,我们实现了持久执行,而没有增加新的攻击面和需要审计的系统。

像 Postgres 这样的数据库已经是许多公司基础设施的一部分,团队已经知道如何大规模运营它们。因此,我们决定重用我们自己的数据库,而不是在其之上硬加上一个外部编排器,这更有意义。如果你已经在运行 Postgres,那么你也已经拥有了持久工作流所需的一切。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存