不想引入 Kafka?在 PostgreSQL 里直接开一个消息队列

发布时间:2026/8/19 19:27:24
不想引入 Kafka?在 PostgreSQL 里直接开一个消息队列 不想引入 Kafka在 PostgreSQL 里直接开一个消息队列【免费下载链接】pgmqA lightweight message queue. Like AWS SQS and RSMQ but on Postgres.项目地址: https://gitcode.com/gh_mirrors/pg/pgmq凌晨两点你的定时任务把一批订单状态同步请求发给了第三方接口接口超时任务进程重启这 500 条消息和它们代表的订单就这样一起丢了。排查到天亮结论是同步脚本缺个可靠的中间环节。这种场景你大概率不陌生——异步任务最怕的就是没地方暂存一崩全丢。PostgreSQL 消息队列PGMQ就是来解决这个问题的它把队列直接建在数据库里不依赖任何外部进程消息随库持久化数据本身就在你眼皮底下。项目定位一句话——像 AWS SQS 和 RSMQ 一样好用的消息队列但跑在 Postgres 上开箱即用。它解决的是什么问题消息队列可以理解成快递中转站生产方把包裹消息放进中转站消费方按自己的节奏来取。没有中转站时生产方和消费方必须掐着点同时在场任何一方掉线业务就卡住。传统做法是引入独立的消息中间件比如 Kafka 或 Redis Streams但这意味着多维护一套集群、多处理一份监控告警、多面对一种网络故障。而 PGMQ 的思路是既然你的业务数据已经在 Postgres 里为什么异步消息不能也住在里面它没有后台工作进程没有额外依赖全部功能就是一组 SQL 函数加几张表。数据库能高可用队列就高可用数据库能备份消息就跟着备份。五分钟把队列跑起来最快路径是用官方 Docker 镜像PGMQ 已作为扩展预装其中docker run -d --name pgmq-postgres -e POSTGRES_PASSWORDpostgres -p 5432:5432 ghcr.io/pgmq/pg18-pgmq:v1.10.0连接数据库后执行一条 SQL 启用扩展CREATE EXTENSION pgmq;如果你想让现有实例直接装上也可以克隆仓库后用 psql 导入 SQL 文件GitHub 镜像源git clone https://gitcode.com/gh_mirrors/pg/pgmq之后psql -f pgmq-extension/sql/pgmq.sql即可。PGMQ 支持 PostgreSQL 14 到 18兼容性覆盖了目前绝大多数生产环境。第一个可运行的完整流程下面是一段完整的创建队列 → 发消息 → 取消息 → 确认处理链路照抄就能跑-- 1. 创建队列本质是在 pgmq schema 下建一张 q_ 前缀的表 SELECT pgmq.create(order_events); -- 2. 发一条 JSON 消息返回消息 id SELECT pgmq.send(order_events, {order_id: A1001, action: paid}); -- 3. 取消息vt30 表示让这条消息对其他消费者隐藏 30 秒 SELECT * FROM pgmq.read(order_events, vt 30, qty 1); -- 4. 处理成功后删除想留底就把 delete 换成 archive SELECT pgmq.delete(order_events, 1);注意read和pop的区别read只是把消息借出来处理完必须手动delete或archivepop则是读取后立即删除适合一次性消费的场景。三个值得深入的能力可见性超时是可靠性的地基。消费者取走消息后它会在指定时间内对其他消费者隐身如果处理进程在这期间崩溃超时一到消息自动恢复可见被下一个消费者重新领取。这等于给恰好一次投递上了保险——既不丢消息也不重复处理。推荐把vt设得比正常处理耗时略长给重试留出余量。FIFO 分组让有序和并发兼得。有些业务要求严格按序比如同一订单的创建、支付、履约另一些又希望不同订单能并行。PGMQ 用消息头里的x-pgmq-group字段做分组同组消息严格先进先出不同组之间互不阻塞可被并行消费。配合pgmq.read_grouped_head()一次取每个组最老的一条消息多进程水平扩展时天然互不抢单。存档比删除更适合审计场景。pgmq.archive()不会物理抹掉消息而是把它挪到a_前缀的存档表里长期保存后续可重放、可对账。对金融交易这类需要留痕的业务存档是默认选项而不是可选项。高流量场景怎么扛消息量大时单表队列会拖慢读写。PGMQ 提供了分区队列pgmq.create_partitioned(high_volume_queue, daily, 30 days)会按时间建分区表配合 pg_partman 自动维护新分区和过期分区。另有pgmq.metrics(my_queue)查看队列积压量、最早消息年龄等指标方便你在堵塞发生前提前干预。相关文档在 docs/partitioned-queues.mdFIFO 细节在 docs/fifo-queues.md接口全量清单在 docs/api/sql/functions.md。几个实操提醒高频使用 FIFO 分组读取的队列建议先执行pgmq.create_fifo_index(my_queue)建索引否则分组查询会退化成全表扫描。消息一旦被读出就进入了处理中状态务必设置合适的vt并及时确认避免积压消息长期占用。丢掉的不是消息而是重试策略——pgmq.set_vt可以把消息的可见性立刻归零强制重试失败到上限再归档形成完整的死信处理闭环。PGMQ 的哲学是少一个组件少一类故障。它把消息队列从一项独立基础设施降维成数据库里的几张表和一组函数——对中小团队这意味着省掉一套 Kafka 集群的运维成本对已经重度使用 Postgres 的团队这意味着异步能力与业务数据天然同源。下次再遇到异步任务怕丢的需求别急着引入新中间件。在你的 PostgreSQL 里先CREATE EXTENSION pgmq;跑通上面那四步 SQL你会很快感受到这套方案的轻巧。【免费下载链接】pgmqA lightweight message queue. Like AWS SQS and RSMQ but on Postgres.项目地址: https://gitcode.com/gh_mirrors/pg/pgmq创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考