别再用状态字段硬撑队列:用 PgQue 在 PostgreSQL 中构建低膨胀工作流

2026-09-28 33 预计阅读时间: 1 分钟
来源: postgr.es AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:13 分钟

很多数据库里的性能问题,并不是查询突然变复杂了,而是某张名为 task、queue 或 bucket 的表每天被更新数百万次:任务从 pending 变成 running,再变成 done 或 failed。随之而来的是反复运行的 autovacuum、持续增长的表膨胀、热点行竞争,以及越来越不稳定的处理延迟。

PgQ 与其现代化实现 PgQue 提供了另一种思路:不要反复更新事件本身,而是把事件写入只追加日志,让每个消费者单独维护自己的读取位置。这种模型仍然完全建立在 PostgreSQL 上,但更接近 Kafka Topic 或 Pulsar Subscription,而不是传统的“取出后删除”队列表。

状态表为何容易进入性能恶性循环

一个常见的任务表大致如下:

create table tasks (
    id bigint generated always as identity primary key,
    status text not null default 'pending',
    payload jsonb not null,
    attempts integer not null default 0,
    updated_at timestamptz not null default now()
);

多个 worker 通常通过下面的模式抢任务:

select id, payload
from tasks
where status = 'pending'
order by id
for update skip locked
limit 100;

FOR UPDATE SKIP LOCKED 本身是非常实用的 PostgreSQL 能力,但它无法消除任务生命周期中的频繁更新。worker 仍要写入 running、done、重试次数和错误信息。PostgreSQL 的 MVCC 会为更新生成新的行版本,旧版本则要等待 vacuum 清理。

当生产速度、消费速度和 vacuum 能力不匹配时,会出现几个相互强化的问题:

  • 高频更新制造大量 dead tuples;
  • autovacuum 频繁扫描热点表和索引;
  • 长事务或落后消费者拖延旧版本清理;
  • 表与索引膨胀导致更多缓存失效和随机 I/O;
  • worker 数量增加后,锁竞争和写放大进一步上升。

这并不意味着 SKIP LOCKED 不可用。对于规模有限、任务保留时间短、消费者模型简单的系统,它往往足够好。问题在于:当一张业务表实际上已经承担事件日志和工作流引擎的职责时,继续优化状态字段只能缓解症状。

PgQue 的关键变化:事件不动,游标前进

PgQue 继承了 PgQ 的核心机制:快照批处理、tick、分区表轮转,以及每个消费者独立维护的游标。

生产者只向分区化的事件日志追加数据,不在原事件上更新 processed = true。ticker 按固定间隔创建 tick,两个 tick 之间已经完成提交的事务组成一个批次。消费者读取批次并确认后,移动的是自己的游标,而不是事件记录。

这带来了几个直接结果:

  1. 队列日志不依赖高频状态更新。 主要写入模式是 append-only,从根源上减少了队列表的更新膨胀。
  2. 多个订阅者可以共享一份事件数据。 每个注册消费者维护独立位置,都能看到同一批事件,无需为订阅者复制消息。
  3. 确认以批次为单位。 游标指向 tick,而不是某一条消息偏移量。消费者确认的是“已经处理到某个 tick”。
  4. 空间回收受最慢消费者约束。 旧分区通过轮转和 TRUNCATE 回收,但必须等待最慢消费者越过对应范围。停滞消费者不会直接丢数据,却会阻塞回收。

PgQue 默认 tick 周期为 100 毫秒,可调整为 1000 毫秒的精确约数。ticker 会在生成 tick 后发送形如 pgque_<queue> 的 NOTIFY,客户端可以使用 LISTEN/NOTIFY 等待新批次,而不是持续轮询。

需要注意,PgQue 并不是把 PostgreSQL 变成了完整的 Kafka 替代品:

  • 它没有 Kafka 式 partition,因此不能直接提供按 key 分区后的顺序保证;
  • cursor 指向 tick,而非单条消息 offset;
  • ack 是批次级操作,不是逐消息确认;
  • 保留窗口与表轮转相关,默认保证的回退窗口是一个轮转周期,摘要所述默认值为两小时;
  • cooperative consumer 可以让多个 worker 共享逻辑消费者,但该能力在所述版本中仍标记为实验性。

因此,PgQue 更适合已经重度依赖 PostgreSQL、希望减少外部组件,同时能够接受 tick 批处理语义的应用。

用 SQL 跑通生产、消费、重试与 DLQ

下面示例假设 PgQue 已安装,orders 队列已经启动,processor 消费者也已完成订阅,后台 ticker 正常运行。具体安装和队列初始化应按照所使用版本的文档执行,不要凭示例猜测管理函数的参数。

生产一条订单事件:

select pgque.send(
    'orders',
    '{"order_id": 43, "total": 10.00}'::jsonb
);

对于低流量队列,事件数可能尚未达到自动生成下一个 tick 的条件。测试时可以显式请求下一个 tick:

select pgque.force_next_tick('orders');

消费者一次读取最多 100 条事件:

select *
from pgque.receive('orders', 'processor', 100);

空批次不需要确认,receive() 会自动越过它。非空批次则必须调用 ack(),否则下一次读取仍会返回尚未确认的批次:

select pgque.ack(8);

真正需要小心的是部分失败。假设批次 8 中的消息 2004 因下游返回 503 而失败,必须先 nack 失败消息,再确认整个批次:

select pgque.nack(
    8,
    m,
    '30 seconds',
    'downstream 503'
)
from pgque.receive('orders', 'processor', 100) as m
where m.msg_id = 2004;

select pgque.ack(8);

执行顺序不能颠倒。ack(8) 会推进消费者游标;失败事件必须在此之前被送入该消费者自己的延迟重试路径。30 秒后,它会重新进入可消费范围。

还可以限制最大重试次数,让持续失败的消息进入死信队列:

select pgque.set_queue_config(
    'orders',
    'max_retries',
    '2'
);

PgQue 同时提供死信检查、重放和清理 API,包括 dlq_inspect()、dlq_replay()、dlq_replay_all() 与 dlq_purge()。生产环境中应把 DLQ 数量、最老死信年龄和重放结果纳入告警,而不是只把 DLQ 当成永久垃圾箱。

一个可改造的 Python 消费循环

如果应用语言没有现成客户端,也可以直接调用 SQL API。下面示例使用 Psycopg 3,并假设消息处理函数具有幂等性。

安装依赖并配置连接串:

python -m pip install 'psycopg[binary]>=3.1'
export DATABASE_URL='postgresql://app:secret@127.0.0.1:5432/appdb'

保存为 consumer.py:

import os
import time
import psycopg
from psycopg.rows import dict_row

DSN = os.environ["DATABASE_URL"]
QUEUE = "orders"
CONSUMER = "processor"
BATCH_SIZE = 100


def handle_message(message: dict) -> None:
    # 替换为真实业务处理。生产环境应使用 msg_id 或业务键做幂等控制。
    print(
        f"processing msg_id={message['msg_id']} "
        f"type={message['type']} payload={message['payload']}"
    )


def main() -> None:
    # autocommit 让 receive 不长期占用事务;nack 与 ack 再放入一个短事务。
    with psycopg.connect(DSN, autocommit=True, row_factory=dict_row) as conn:
        while True:
            messages = conn.execute(
                "select * from pgque.receive(%s, %s, %s)",
                (QUEUE, CONSUMER, BATCH_SIZE),
            ).fetchall()

            if not messages:
                time.sleep(0.5)
                continue

            batch_id = messages[0]["batch_id"]
            failures: list[tuple[int, str]] = []

            for message in messages:
                try:
                    handle_message(message)
                except Exception as exc:
                    failures.append((message["msg_id"], str(exc)))

            # 所有 nack 必须发生在 ack 之前,并尽量在同一个短事务中完成。
            with conn.transaction():
                for msg_id, reason in failures:
                    conn.execute(
                        """
                        select pgque.nack(
                            %s,
                            m,
                            '30 seconds',
                            %s
                        )
                        from pgque.receive(%s, %s, %s) as m
                        where m.msg_id = %s
                        """,
                        (
                            batch_id,
                            reason[:500],
                            QUEUE,
                            CONSUMER,
                            BATCH_SIZE,
                            msg_id,
                        ),
                    ).fetchone()

                conn.execute(
                    "select pgque.ack(%s)",
                    (batch_id,),
                ).fetchone()


if __name__ == "__main__":
    main()

运行:

python consumer.py

示例故意没有把业务处理和数据库确认放进一个长事务。这样可以降低长事务对 PostgreSQL 的影响,但也意味着进程可能在“外部操作成功、批次尚未 ack”之间崩溃,随后收到重复事件。因此消费者仍应按至少一次投递来设计:使用 msg_id、订单号或业务幂等键记录处理结果,避免重复扣款、重复发货或重复调用外部 API。

生产者同样可能在“写入成功、客户端未收到响应”时重试。PgQue 提供 send_idem() 和相关维护能力用于幂等发布,适合有稳定业务去重键的场景。

上线前不要只看吞吐量

决定采用 PgQue 前,可以按下面的清单评估:

  • 消息顺序是否依赖 key partition? 如果需要海量分区、按 key 严格有序,Kafka 一类系统通常更匹配。
  • 消费者能否接受批次 ack? 部分失败必须先逐条 nack,再 ack 整批。
  • 处理逻辑是否幂等? 崩溃、超时和重试都可能带来重复投递。
  • 能否及时发现停滞消费者? 最慢消费者会阻塞旧表回收,应监控 stuck_consumers()、消费者延迟和磁盘增长。
  • 数据库是否还有足够隔离空间? 工作流与核心 OLTP 共用 PostgreSQL 可以减少组件数量,也可能让流量尖峰争夺同一组 CPU、I/O 和连接。
  • 权限是否最小化? PgQue 区分 pgque_writer、pgque_reader 和 pgque_admin,应用账号不应默认获得管理权限。
  • 是否建立 DLQ 处置流程? 进入死信队列只是隔离故障,不等于故障已经解决。

PgQue 最有价值的地方,并不是一句“用 PostgreSQL 替代 Kafka”,而是把数据库内的工作流从高频状态更新改造成追加日志与消费者游标。对于中等规模、强事务关联、运维团队希望减少基础设施种类的系统,这种设计能显著降低队列表膨胀和 vacuum 压力;而当需求扩展到跨区域日志、超大吞吐、复杂分区和长期独立保留策略时,专用消息平台仍然有清晰优势。


相关推荐