
Postgres 存数据、跑查询都没问题,但一旦需要跑多步骤任务——夜间聚合、生成 embedding、等人工审,你就得搭一整套外部编排系统。队列表、轮询 worker、状态机、崩溃恢复……本质上就是为了补一个 Postgres 缺失的原语:持久化的异步后台执行。微软开源的 pg_durable 扩展,就是来填这个坑的。
大多数 PG 团队跑后台任务,最后都走上两条路:
路径一:塞进 PL/pgSQL 函数 一个连接从头占到尾,所有逻辑跑在一个大事务里。连接断了?全丢。数据库重启?全丢。想重试某个步骤?没门。想并行?没门。想暂停等人审批?更没门。
路径二:搬到外部 任务队列、轮询 worker、状态表、重试逻辑、崩溃清理……几个后台任务最终变成一个完整的分布式系统。而你要编排的工作,本质上根本没离开过数据层。 两条路都在给同一个缺失的原语打补丁:持久化、异步的后台工作,存在于数据所在的位置。
pg_durable 由两部分组成:
组件 | 说明 |
|---|---|
DSL(领域特定语言) | 用 SQL 表达式描述工作流,支持顺序、并行、条件、循环 |
duroxide 运行时 | 托管在 Postgres 后台 worker 中,异步执行 |
核心机制:调用 df.start(...) 立即返回一个实例 ID,工作在后台异步跑。每个步骤在独立会话和事务中执行,提交进度后移交下一步。数据库崩溃?从检查点恢复。某个步骤失败?只重试那一步,已成功的不重跑。
运算符 | 含义 | 类比 |
|---|---|---|
~> | 顺序执行 | 先 A 再 B |
& | 并行执行 | 扇出,等全部完成 |
| | 竞速执行 | 扇出,取第一个完成的 |
?> / !> | 条件分支 | if / else |
@> | 持久化循环 | 重启后仍能存活 |
|=> | 结果捕获到变量 | 赋值给 $var |
函数 | 用途 |
|---|---|
df.start() | 启动工作流,返回实例 ID |
df.status() | 查看执行状态 |
df.result() | 获取执行结果 |
df.http() | 调用允许列表中的 HTTP 端点 |
df.wait_for_schedule() | Cron 风格定时调度 |
df.wait_for_signal() | 暂停等待外部事件(人工审批) |
-- 直接创建扩展
CREATE EXTENSION IF NOT EXISTS pg_durable;# 克隆仓库
git clone https://github.com/microsoft/pg_durable.git
# 使用 Codespace 预构建或 VS Code Dev Container
# 要求:PostgreSQL 17+pg_durable GitHub 仓库:https://github.com/microsoft/pg_durable
-- 启动一个持久化工作流
SELECT df.start($$ SELECT 'Hello, durable world!' AS message $$);
-- 返回:a1b2c3d4(8 字符实例 ID)
-- 查看结果
SELECT df.result('a1b2c3d4');
-- message: Hello, durable world!提交、离开、回来检查。不阻塞连接,不占用事务。
SELECT df.start(
'DELETE FROM target WHERE loaded_at < now() - interval ''7 days'''
~> 'UPDATE staging SET processed_at = now() WHERE processed_at IS NULL'
~> 'INSERT INTO target (data, source_id)
SELECT data, source_id FROM staging WHERE processed_at IS NOT NULL',
'etl-pipeline'
);如果数据库在第三步跑到一半重启了,从最后一个检查点批次继续,而不是从第零行重来。
SELECT df.start(
@> (
(df.http('https://httpbingo.org/json', 'GET') |=> 'response')
~> 'INSERT INTO external_data_sync (data) VALUES ($response::jsonb)'
~> df.wait_for_schedule('*/30 * * * *') -- 每 30 分钟
),
'scheduled-data-sync'
);@> 创建的循环在数据库重启后仍然存活。
SELECT df.start(
'SELECT amount > 10000 AS needs_review FROM invoices WHERE id = 42' |=> 'risky'
?> ( df.wait_for_signal('invoice-42')
~> 'UPDATE invoices SET status = ''paid'' WHERE id = 42' )
!> 'UPDATE invoices SET status = ''paid'' WHERE id = 42',
'invoice-approval'
);金额超过 1 万的发票,暂停等待审核人员发出信号。可以等几分钟,也可以等几天。
同样的需求:"并行跑三个聚合,然后刷新物化视图,带重试和崩溃恢复"。
手写版:
组件 | 工作量 |
|---|---|
队列表 | 1 张 |
轮询 worker | 1 个 |
状态机行 | N 行 |
按步骤重试逻辑 | 自行实现 |
崩溃恢复清理 | 自行实现 |
总代码量 | 300+ 行 |
pg_durable 版:
SELECT df.start(
'SELECT count(*) FROM users'
& 'SELECT count(*) FROM orders'
& 'SELECT sum(amount) FROM orders'
~> 'REFRESH MATERIALIZED VIEW metrics',
'refresh-dashboard'
);你写 SQL,pg_durable 管队列、状态、协调、重试和崩溃恢复。
在 HorizonDB 上,pg_durable 驱动 azure_ai 扩展的声明式 AI 管道:
SELECT ai.create_pipeline(
name => 'ai_pipeline',
source => ai.table_source(table_name => 'documents'),
steps => ARRAY[
ai.chunk(input => 'content'),
ai.embed(model => 'default-embedding', input => 'chunk_text', dimensions => 1536)
],
trigger => 'on_change',
sink => ai.table_sink('documents_output')
);
SELECT ai.run('ai_pipeline');每个 AI 步骤是一个持久化节点。ai.embed() 失败了?ai.chunk() 不会重跑。trigger => 'on_change' 只嵌入新内容。加上 DiskANN 索引,端到端向量搜索全在数据库里。
夜间数据清理、周期性统计刷新、定时数据归档。用 df.wait_for_schedule() + @> 循环,重启不丢。
需要人工审批的订单处理、跨步骤数据转换、条件分支的数据管道。用 df.wait_for_signal() 暂停,用 ?> / !> 分支。
pg_durable 不是通用编排器。
场景 | 是否适合 | 原因 |
|---|---|---|
工作流与 PG 数据紧耦合 | ✅ 适合 | 读写同一数据库,享受 PG 持久性和备份 |
跨异构服务扇出 | ❌ 不适合 | 用专用编排器(如 Temporal) |
运行任意应用逻辑 | ❌ 不适合 | 无法映射为 SQL 步骤/分支/循环/HTTP 调用 |
ETL / 嵌入 / 定时维护 | ✅ 适合 | 数据不离开数据层 |
PG 17 以下版本 | ❌ 不支持 | 最低要求 PostgreSQL 17 |
指标 | 数据 |
|---|---|
发布当天 | Hacker News 热门 |
GitHub Stars | 发布几天内 1700+ |
独立教程 | Franck Pachot 发表入门教程 |
许可证 | PostgreSQL License(宽松开源) |
仓库活跃度 | 维护者阅读每个 issue 和 PR |
安装 PostgreSQL 扩展后可以:
pg_durable 的 HTTP 调用只允许预先配置的端点。 出于安全考虑,
df.http()不能调用任意 URL,必须在允许列表中配置。生产环境使用前务必检查网络策略。
pg_durable 解决的问题很明确:把持久化异步执行搬回数据库里。
df.start() 搞定如果你手里有跑在 PG 上的 ETL、定时任务、embedding 管道,强烈建议试一下。数据不离开数据库,编排逻辑全在 SQL 里,运维负担直接砍半。