pg_durable 实战:把可恢复的数据处理流程留在 PostgreSQL 里
批量入库、向量化、数据清洗这类任务,往往从一段定时脚本开始:表里加状态列,应用层维护重试次数,再配队列、worker 和监控。真正麻烦的不是第一次跑通,而是数据库重启、某一批 API 调用失败或任务跑到一半时,怎样确定哪些步骤已经完成、哪些步骤该安全重试。
微软开源的 pg_durable 提供了另一条路径:把长时间运行、可恢复的工作流定义为 PostgreSQL 内的 SQL 步骤。它是用 Rust/pgrx 构建的 PostgreSQL 扩展,采用 PostgreSQL License;运行时在数据库内保存流程状态和检查点。发生崩溃、重启或步骤失败后,已完成步骤无需重新执行,未完成部分可从最近的持久化检查点继续。
这不是要替代所有编排平台。它特别适合工作状态本来就在 Postgres 中、希望少维护一层队列和状态表的批处理;若流程主要跨越许多异构系统、依赖任意 SDK 或复杂的内存控制流,独立的工作流引擎仍可能更合适。
先理解它解决的边界
pg_durable 将一个任务表示为 SQL 步骤图:df.start() 提交任务,后台 worker 逐步执行,并在步骤之间记录进度。顺序执行可用 ~> 连接;也可以使用并行、条件、定时等待和信号等原语。运行实例和节点状态可从 PostgreSQL 中查询,不必另外拼接 worker 日志、队列表和状态同步逻辑。
例如,一个文档处理流程通常要取未处理记录、分批标记,再写回结果。下面的示例只展示最小的两步持久化流程,SQL 与 df.start() 的调用方式来自项目 README:
SELECT df.start(
'SELECT id FROM documents WHERE processed = false LIMIT 100' |=> 'batch'
~> 'UPDATE documents
SET processed = true
WHERE id IN (SELECT id FROM $batch.*)'
);
|=> 'batch' 为前一步结果命名,后一段通过 $batch 引用它。实际生产任务应把「已处理」替换成具有业务语义的幂等写入:例如写入带唯一键的处理结果,或在同一事务语义下更新状态。可恢复不等于自动消除副作用;如果某一步会调用外部 HTTP 服务,仍应在业务接口侧设计幂等键和重试策略。
安装时有三个不可省略的动作
项目为 PostgreSQL 17 与 18 发布 Debian 包和 GHCR 镜像,也支持从源代码构建。官方 README 明确将已发布 Docker 镜像定位为评估和学习用途,而不是生产部署方案;生产环境应按项目的包、权限与运维文档完成评估。
无论采用哪种安装方式,扩展的后台 worker 都需要在 PostgreSQL 启动时加载。先在 postgresql.conf 的 shared_preload_libraries 中加入 pg_durable,重启数据库,再创建扩展:
CREATE EXTENSION pg_durable;
-- 将使用权限授予应用角色,而不是开放给所有用户
SELECT df.grant_usage('app_role');
创建扩展后,后台 worker 会异步初始化内部 schema,短暂等待后再调用 df.* 函数即可。还应注意两个权限事实:CREATE EXTENSION 不会向 PUBLIC 自动授权;worker 角色需要 superuser 权限以管理所有用户的实例。因此,把业务连接限制在普通应用角色、只授予 df.grant_usage() 所需权限,比分发超级用户连接串更稳妥。
从实例 ID 开始建立可观测性
提交任务会返回实例 ID。将该 ID 写入业务日志或关联表,就能在故障排查时回到数据库中查看任务状态和结果:
-- 提交一个可恢复任务
SELECT df.start('SELECT ''Hello, durable world!''');
-- 使用实际返回的实例 ID 查询状态与结果
SELECT df.status('a1b2c3d4');
SELECT df.result('a1b2c3d4');
这里的 a1b2c3d4 只是示例 ID。对批处理而言,更实用的做法是在提交后立刻保存真实 ID,并把取消、重试和告警围绕实例状态实现。项目也提供 df.list_instances()、df.cancel()、df.explain() 等接口,适合把任务进度接入现有的 DBA 或 SRE 查询流程。
适合落地的三个场景
第一,向量或文档入库:按批读取记录、调用已做幂等保护的服务、再更新处理状态。第二,数据维护:检测膨胀、等待人工确认、执行下一步维护动作。第三,聚合类作业:将多个独立查询并行执行,最后合并为报表结果。它们的共同点是数据、状态与审计需求都紧邻 PostgreSQL。
开始试用前,建议先用一个小而可重复的任务验证四件事:扩展重启后的恢复行为、步骤的幂等性、应用角色权限,以及实例状态如何进入现有监控。若这些边界清楚,pg_durable 能让一部分原本散落在 cron、队列、worker 和状态表之间的流程,回到最接近数据的地方运行。
相关链接