pgmq

基于Postgres实现类似AWS SQS/RSMQ的消息队列

概览

扩展包名版本分类许可证语言
pgmq1.12.0FEATPostgreSQLSQL
ID扩展名BinLibLoadCreateTrustReloc模式
2660pgmqpgmq
相关扩展kafka_fdw pg_task pg_net pg_background pgagent pg_jobmon
下游依赖fsm_core pg_later vectorize

版本

类型仓库版本PG 大版本包名依赖
EXTPIGSTY1.12.01817161514pgmq-
RPMPIGSTY1.12.01817161514pgmq_$v-
DEBPIGSTY1.12.01817161514postgresql-$v-pgmq-
OS / PGPG18PG17PG16PG15PG14
el8.x86_64
el8.aarch64
el9.x86_64
el9.aarch64
el10.x86_64
el10.aarch64
d12.x86_64
d12.aarch64
d13.x86_64
d13.aarch64
u22.x86_64
u22.aarch64
u24.x86_64
u24.aarch64
u26.x86_64
u26.aarch64

构建

您可以使用 pig build 命令构建 pgmq 扩展的 RPM / DEB 包:

pig build pkg pgmq         # 构建 RPM / DEB 包

安装

您可以直接安装 pgmq 扩展包的预置二进制包,首先确保 PGDGPIGSTY 仓库已经添加并启用:

pig repo add pgsql -u          # 添加仓库并更新缓存

使用 pig 或者是 apt/yum/dnf 安装扩展:

pig install pgmq;          # 当前活跃 PG 版本安装
pig ext install -y pgmq -v 18  # PG 18
pig ext install -y pgmq -v 17  # PG 17
pig ext install -y pgmq -v 16  # PG 16
pig ext install -y pgmq -v 15  # PG 15
pig ext install -y pgmq -v 14  # PG 14
dnf install -y pgmq_18       # PG 18
dnf install -y pgmq_17       # PG 17
dnf install -y pgmq_16       # PG 16
dnf install -y pgmq_15       # PG 15
dnf install -y pgmq_14       # PG 14
apt install -y postgresql-18-pgmq   # PG 18
apt install -y postgresql-17-pgmq   # PG 17
apt install -y postgresql-16-pgmq   # PG 16
apt install -y postgresql-15-pgmq   # PG 15
apt install -y postgresql-14-pgmq   # PG 14

创建扩展

CREATE EXTENSION pgmq;

用法

来源:

pgmq 实现了持久化消息队列,作为 PostgreSQL 表和 SQL 函数。它支持延迟投递、可见性超时、FIFO 组、消息头、轮询、主题和归档功能。当需要将队列事务与同一数据库中的关系变化协调起来时,请使用此扩展。

创建队列并发送消息

CREATE EXTENSION pgmq;
SELECT pgmq.create('jobs');

SELECT *
FROM pgmq.send(
  queue_name => 'jobs',
  msg        => '{"task":"refresh"}'::jsonb,
  delay      => 0
);

send 返回消息标识符。send_batch 可插入多个 JSONB 消息。头部可以携带路由或跟踪元数据,这些元数据与主体分开存储,并且在某些重载中支持它们。

使用可见性超时读取消息

SELECT *
FROM pgmq.read(
  queue_name => 'jobs',
  vt         => 30,
  qty        => 10
);

读取操作会隐藏每条消息 vt 秒。成功后,可以删除或归档它:

SELECT pgmq.delete('jobs', 42);
SELECT pgmq.archive('jobs', 43);

如果处理失败或者消费者消失,则未确认的消息将再次可见。因此,消费者必须是幂等的;pgmq 不会保证任意应用程序副作用在全局范围内恰好执行一次。

pop 读取消息并立即删除,仅当允许在调用后丢失消息时才适用。

FIFO 组头轮询

1.12.0 版本增加了对多个 FIFO 组头部的消息轮询:

SELECT *
FROM pgmq.read_grouped_head_with_poll(
  queue_name          => 'jobs',
  vt                  => 30,
  qty                 => 10,
  max_poll_seconds    => 5,
  poll_interval_ms    => 100
);

此操作会选择组头部消息并进行轮询,直到达到最大轮询时间。这保留了每个组内的顺序,并允许不同组并发处理。

队列管理索引

  • pgmq.create(queue_name): 创建队列和归档结构。
  • pgmq.send 和 pgmq.send_batch: 入队 JSONB 消息,可选延迟。
  • pgmq.read: 为可见性超时声明消息。
  • pgmq.read_grouped_head_with_poll: 轮询 FIFO 组头部。
  • pgmq.pop: 读取消息并立即删除。
  • pgmq.delete: 通过移除消息来确认。
  • pgmq.archive: 将消息移动到队列归档表中。
  • pgmq.drop_queue: 移除队列对象。
  • pgmq.metrics 和相关辅助函数: 检查可用时的队列深度和年龄。

对于队列作业,归档行存储在 pgmq.a_<queue_name> 中。将这些表视为由扩展管理的对象。

运营注意事项

  • 将 vt 设定为比正常处理时间更长,并设计好超时后的重新投递。
  • 队列和归档表消耗普通 PostgreSQL 的 WAL、存储、真空和备份容量。
  • 归档或删除已完成的消息并强制执行归档保留策略。
  • 长期轮询会占用数据库连接。根据消费者数量调整连接池大小和轮询间隔。
  • 保持队列名称符合 pgmq 的标识符规则;调用 API 而不是从不可信输入构造表名。