pg_durable
microsoft/pg_durable加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
加载项目详情…
本应用为开源项目,仅供学习研究,请遵守其开源协议。
一句话定位:Microsoft 出品的 PostgreSQL 17+ 扩展,用纯 SQL 定义工作流,无需外部服务即可实现故障恢复、自动重试与并行执行。
凌晨三点,值班的 DBA 收到一条告警:数据库主库意外重启。在此之前,一条 ETL 流水线跑了 6 个小时,正在执行"第三步:将数据写入聚合表"——这条 SQL 跑完后,理论上会更新 pipeline_status 标记,进入下一步。
但数据库在这个精确的瞬间重启了。
重启后的系统看到的是:第三步"似乎完成了"(进程崩溃前提交了事务),但 pipeline_status 还没更新。流水线实际上是未完成状态。监控误报"一切正常",直到下一个批次触发唯一键冲突,才发现数据重复了 30%。
这就是没有持久化执行的代价:应用层的"任务进度"和数据库实际状态之间存在一道裂缝,一旦进程崩溃,这道裂缝就会吞噬数据一致性。
pg_durable 要解决的问题,就是彻底消除这道裂缝。
pg_durable 由 Microsoft 内部团队开发,核心目标在官方文档中写得很清楚:"bring compute close to data"——让计算靠近数据本身。
这个思路并不新鲜。Microsoft 的 Durable Functions(Azure Functions 的扩展)在云端早已是成熟范式,其底层引擎 Duroxide 是一个 Rust 写的持久化执行运行时。pg_durable 本质上是把 Duroxide 移植到了 PostgreSQL 内部,用 PostgreSQL 本身作为持久化存储,去掉了对 Temporal 等外部协调服务的依赖。
2026 年中旬,Microsoft 将 pg_durable 完全开源,同时将其内置于 Azure HorizonDB(微软全新 PostgreSQL 云服务)中。HorizonDB 的 AI 流水线层也构建在 pg_durable 之上,每一阶段都经过 checkpoint、失败重试保障,从原始数据到可用向量的全流程 crash-safe。
pg_durable 的设计哲学是**"工作流即 SQL,SQL 即工作流"**。它提供了一套 SQL 领域的 DSL(领域特定语言),用熟悉的操作符组合出复杂工作流。
| 操作符 | 含义 | 类比 |
|---|---|---|
~> | 顺序执行(sequential) | 流水线 Pipe |
& | 并行执行(fan-out) | Promise.all / asyncio.gather |
| ` | =>` | 变量绑定(assign) |
-- 安装扩展
CREATE EXTENSION pg_durable;
-- 最简单的示例:输出一句话
SELECT df.start('SELECT ''Hello, durable world!'' AS message');
-- 一个并行聚合:三个查询同时跑,结果汇总到 dashboard
SELECT df.start(
'SELECT count(*) FROM users' &
'SELECT count(*) FROM orders' &
'SELECT sum(amount) FROM orders'
~> 'REFRESH MATERIALIZED VIEW dashboard',
'metrics'
);
df.start() 立即返回一个 instance_id(如 a1b2c3d4),后台 worker 开始异步执行。开发者无需等待,随时可以用 df.explain() 查看执行状态:
SELECT df.explain('a1b2c3d4');
-- Instance: a1b2c3d4
-- Status: ✓ Completed
-- Output: {"rows": [...], "row_count": 1}
| 场景 | 传统方案痛点 | pg_durable 优势 |
|---|---|---|
| 向量嵌入流水线 | 进程中断导致 embedding 重复或丢失 | 每个步骤 checkpoint,crash 后从最后一步恢复 |
| ETL 管道 | 状态表 + 轮询 worker,代码复杂 | 顺序链式执行,SQL 可读,天然事务一致性 |
| 定时任务 | pg_cron + 状态列 + 重试逻辑 | 内置 cron 调度 + 自动重试 + 持久化 |
| 人工审批流 | 独立任务队列 + webhook | 条件等待 + 信号机制,SQL 内可表达审批逻辑 |
| 外部 API 调用 | 独立消息队列 + 回调服务 | df.http() 直接从 SQL 发起 HTTP 请求,SSRF 保护内置 |
pg_durable 是一个 PostgreSQL 扩展(Extension),完整运行在 PostgreSQL 进程内部,不依赖任何外部服务(只要有 PostgreSQL 17+ 即可)。
┌─────────────────────────────────────────────┐
│ User Session (Phase 1) │
│ df.start() → 构建有向无环图(DAG) │
│ 操作符: ~>, &=, |=> │
│ 输出: instance_id │
└────────────────┬────────────────────────────┘
│ instance_id 入队
▼
┌─────────────────────────────────────────────┐
│ Background Worker (Phase 2) │
│ 异步执行有向图,checkpoint 每一步 │
│ 底层: Duroxide 持久化执行引擎 │
│ 存储: duroxide schema(PostgreSQL 内) │
└─────────────────────────────────────────────┘
Phase 1(用户会话内,同步):操作符调用 DSL 函数,在用户事务中构建节点图,最后 df.start() 将 instance_id 写入 duroxide 队列。立即返回,用户无需等待。
Phase 2(后台 Worker,异步):独立后台进程从队列取出 instance_id,按拓扑顺序执行各节点,每步完成后立即 checkpoint(写入 PostgreSQL 表)。如果数据库崩溃,Worker 重启后会从最后一个 checkpoint 恢复。
| 库 | 版本 | 作用 |
|---|---|---|
| pgrx | 0.16.1 | PostgreSQL 扩展开发框架(Rust ↔ Postgres FFI) |
| duroxide | 0.1.30 | 持久化执行引擎,checkpoint + 重放核心 |
| sqlx | 0.8 | ExecuteSQL 活动类型(Worker 重新连接 Postgres 执行 SQL) |
| reqwest | 0.13 | df.http() HTTP 请求(SSRF 保护内置) |
| tokio | 1 | 异步运行时(后台 Worker) |
df.join() 等待所有分支完成并汇合结果|=> 绑定变量,后续步骤用 $变量名 引用wait_for_schedule() 内置 cron 表达式解析df.http() 从 SQL 内发起外部 API 调用continue_as_new 支持长时间循环工作流开发者友好度:中等。官方提供了 Docker 一键启动方案(docker-compose up),适合快速试用:
git clone https://github.com/microsoft/pg_durable
cd pg_durable
docker compose up
docker exec -it pg_durable psql -U postgres -d pg_durable
原生编译需要:Rust toolchain + PostgreSQL 17 开发包 + pgrx 工具链,约 5-10 分钟。官方 Multi-stage Dockerfile 对编译步骤做了完整封装。
pg_durable 的 DSL 非常符合直觉——熟悉 SQL 的开发者无需学习新语言。例如,构建一个文档处理流水线:
SELECT df.start(
'SELECT id FROM documents WHERE processed = false LIMIT 100' |=>
'batch'
~> 'SELECT embed($batch) FROM documents WHERE id = ANY($batch)' |=>
'embeddings'
~> 'UPDATE documents SET processed = true WHERE id = ANY($batch)'
);
整段代码无需外部配置,无需定义状态表,逻辑一目了然。
所有执行状态都在 PostgreSQL 表中,通过标准 SQL 查询即可监控:
-- 列出所有工作流实例
SELECT * FROM df.instances;
-- 查看某个实例的执行详情
SELECT df.explain('instance_id');
-- 查看历史记录
SELECT * FROM df.instances WHERE status = 'completed';
对于 Azure HorizonDB 用户,还可以通过 VS Code 的 PostgreSQL 扩展直接在编辑器内查看工作流状态。
目前只支持 PostgreSQL 17 和 18。大量企业仍在使用 PG 14/15/16,升级数据库主版本并非易事。这意味着 pg_durable 的采纳速度在很大程度上受制于企业的数据库升级周期。
PostgreSQL 扩展的编译比 Python pip install 或 Node npm install 复杂得多。需要 Rust 工具链、PostgreSQL 开发头文件,跨平台编译(Windows/macOS)需要额外配置。Docker 方案降低了门槛,但对习惯了纯 SQL 运维的 DBA 来说仍有学习曲线。
Temporal、Airflow、Dagster 等外部编排工具拥有更丰富的生态(可视化 DAG、任务重试策略、分布式执行、与其他数据源的深度集成)。pg_durable 的优势在于靠近数据(无需跨网络传输数据),但对于跨多个异构数据源的工作流,外部编排工具可能更合适。
2026 年,数据库内工作流编排正在成为一条热门赛道。DBOS 的 pg-queue、surfql、以及 pg_durable 先后开源,都在探索同一个命题:与其在应用层维护复杂的状态机,为什么不把工作流交给数据库?
pg_durable 的差异化在于:
CREATE EXTENSION 即可
适合的场景:
不太适合的场景:
一句话评价:pg_durable 把"持久化执行"这个过去需要整套分布式系统才能实现的能力,压缩进了一个 PostgreSQL 扩展。如果你的工作流主要跟 PostgreSQL 打交道,这个扩展值得认真评估。