首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >让 PostgreSQL 自己跑后台任务:pg_durable 扩展把队列和重试逻辑全干掉了

让 PostgreSQL 自己跑后台任务:pg_durable 扩展把队列和重试逻辑全干掉了

作者头像
小徐
发布2026-07-24 21:16:19
发布2026-07-24 21:16:19
1320
举报
文章被收录于专栏:GreenplumGreenplum

让 PostgreSQL 自己跑后台任务:pg_durable 扩展把队列和重试逻辑全干掉了

Postgres 存数据、跑查询都没问题,但一旦需要跑多步骤任务——夜间聚合、生成 embedding、等人工审,你就得搭一整套外部编排系统。队列表、轮询 worker、状态机、崩溃恢复……本质上就是为了补一个 Postgres 缺失的原语:持久化的异步后台执行。微软开源的 pg_durable 扩展,就是来填这个坑的。

两条路都是死胡同

大多数 PG 团队跑后台任务,最后都走上两条路:

路径一:塞进 PL/pgSQL 函数 一个连接从头占到尾,所有逻辑跑在一个大事务里。连接断了?全丢。数据库重启?全丢。想重试某个步骤?没门。想并行?没门。想暂停等人审批?更没门。

路径二:搬到外部 任务队列、轮询 worker、状态表、重试逻辑、崩溃清理……几个后台任务最终变成一个完整的分布式系统。而你要编排的工作,本质上根本没离开过数据层。 两条路都在给同一个缺失的原语打补丁:持久化、异步的后台工作,存在于数据所在的位置。

扩展原理:DSL + 后台 worker

pg_durable 由两部分组成:

组件

说明

DSL(领域特定语言)

用 SQL 表达式描述工作流,支持顺序、并行、条件、循环

duroxide 运行时

托管在 Postgres 后台 worker 中,异步执行

核心机制:调用 df.start(...) 立即返回一个实例 ID,工作在后台异步跑。每个步骤在独立会话和事务中执行,提交进度后移交下一步。数据库崩溃?从检查点恢复。某个步骤失败?只重试那一步,已成功的不重跑。

DSL 运算符速查

运算符

含义

类比

~>

顺序执行

先 A 再 B

&

并行执行

扇出,等全部完成

|

竞速执行

扇出,取第一个完成的

?> / !>

条件分支

if / else

@>

持久化循环

重启后仍能存活

|=>

结果捕获到变量

赋值给 $var

核心函数

函数

用途

df.start()

启动工作流,返回实例 ID

df.status()

查看执行状态

df.result()

获取执行结果

df.http()

调用允许列表中的 HTTP 端点

df.wait_for_schedule()

Cron 风格定时调度

df.wait_for_signal()

暂停等待外部事件(人工审批)

安装

在 Azure HorizonDB 上

代码语言:javascript
复制
-- 直接创建扩展
CREATE EXTENSION IF NOT EXISTS pg_durable;

在本地或任何 Postgres 17 上

代码语言:javascript
复制
# 克隆仓库
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

使用方法

Hello World

代码语言:javascript
复制
-- 启动一个持久化工作流
SELECT df.start($$ SELECT 'Hello, durable world!' AS message $$);
-- 返回:a1b2c3d4(8 字符实例 ID)

-- 查看结果
SELECT df.result('a1b2c3d4');
-- message: Hello, durable world!

提交、离开、回来检查。不阻塞连接,不占用事务。

ETL 管道

代码语言:javascript
复制
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'
);

如果数据库在第三步跑到一半重启了,从最后一个检查点批次继续,而不是从第零行重来。

定时数据同步(永久运行)

代码语言:javascript
复制
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'
);

@> 创建的循环在数据库重启后仍然存活。

人工审批流程

代码语言:javascript
复制
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 万的发票,暂停等待审核人员发出信号。可以等几分钟,也可以等几天。

效果对比:300 行代码 vs 一行 SQL

同样的需求:"并行跑三个聚合,然后刷新物化视图,带重试和崩溃恢复"。

手写版:

组件

工作量

队列表

1 张

轮询 worker

1 个

状态机行

N 行

按步骤重试逻辑

自行实现

崩溃恢复清理

自行实现

总代码量

300+ 行

pg_durable 版:

代码语言:javascript
复制
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 管队列、状态、协调、重试和崩溃恢复。

适用场景

场景一:Embedding 管道(HorizonDB 专属)

在 HorizonDB 上,pg_durable 驱动 azure_ai 扩展的声明式 AI 管道:

代码语言:javascript
复制
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

VS Code 支持

安装 PostgreSQL 扩展后可以:

  1. 1. 直连 HorizonDB 或本地 Postgres
  2. 2. 让 Copilot 用 pg-durable-sql skill 把自然语言转成 DSL 语法("每天晚上归档超过 90 天的订单")
  3. 3. 实时可视化工作流图——每步耗时、失败位置一目了然

兼容性

  • PostgreSQL 17+
  • Azure HorizonDB(AI 管道功能专属)
  • • 任何云、本地服务器、笔记本电脑均可运行
  • • 开源协议:PostgreSQL License

重要提醒

pg_durable 的 HTTP 调用只允许预先配置的端点。 出于安全考虑,df.http() 不能调用任意 URL,必须在允许列表中配置。生产环境使用前务必检查网络策略。

总结

pg_durable 解决的问题很明确:把持久化异步执行搬回数据库里

  • • 原来要 300 行代码搭的队列+重试+崩溃恢复,现在一行 df.start() 搞定
  • • 每个步骤独立事务、独立检查点,崩溃后从断点恢复
  • • AI 管道场景下,chunk 失败不重跑 embed,embed 失败不重跑 chunk
  • • 开源、宽松许可、GitHub 1700+ Stars,社区活跃

如果你手里有跑在 PG 上的 ETL、定时任务、embedding 管道,强烈建议试一下。数据不离开数据库,编排逻辑全在 SQL 里,运维负担直接砍半。

引用链接

  • pg_durable GitHub 仓库:https://github.com/microsoft/pg_durable
  • 原文:Introducing Durable Functions in PostgreSQL:https://techcommunity.microsoft.com/blog/adforpostgresql/introducing-durable-functions-in-postgresql/4526821
  • HorizonDB 持久化函数文档:https://learn.microsoft.com/en-us/azure/horizondb/development/durable-functions
  • HorizonDB AI 管道文档:https://learn.microsoft.com/en-us/azure/horizondb/ai/ai-pipelines
  • Franck Pachot 入门教程:https://dev.to/franckpachot/getting-started-with-pgdurable-durable-workflows-inside-postgresql-3980
本文参与 腾讯云自媒体同步曝光计划,分享自微信公众号。
原始发表:2026-07-23,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 河马coding 微信公众号,前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体同步曝光计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 让 PostgreSQL 自己跑后台任务:pg_durable 扩展把队列和重试逻辑全干掉了
    • 两条路都是死胡同
    • 扩展原理:DSL + 后台 worker
      • DSL 运算符速查
      • 核心函数
    • 安装
      • 在 Azure HorizonDB 上
      • 在本地或任何 Postgres 17 上
    • 使用方法
      • Hello World
      • ETL 管道
      • 定时数据同步(永久运行)
      • 人工审批流程
    • 效果对比:300 行代码 vs 一行 SQL
    • 适用场景
      • 场景一:Embedding 管道(HorizonDB 专属)
      • 场景二:定时维护任务
      • 场景三:多步骤业务流程
    • 限制说明
    • 社区反响
      • VS Code 支持
      • 兼容性
      • 重要提醒
    • 总结
    • 引用链接
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档