用 PgQ 与 PgQue 把 PostgreSQL 高频更新改造成可控工作流

2026-09-28 20 预计阅读时间: 1 分钟
来源: percona.com 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.

预计阅读时间:11 分钟

生产环境里,有一类表看起来不大,却能把 PostgreSQL 拖得很累:同一批记录每天被更新数百万次,autovacuum 一天运行数百次,表和索引仍不断膨胀。最终表现为查询变慢、并发下降,维护窗口也越来越难安排。

PgQ 和 PgQue 代表了另一种设计方向:不要让业务进度持续覆盖同一行,而是把需要执行的动作建模为队列消息或工作流任务。它们并不会自动消除 MVCC、VACUUM 或锁竞争,但能帮助系统把高频状态变更从核心业务表中隔离出来,并建立清晰的领取、重试和完成机制。

为什么反复 UPDATE 会变成数据库负担

PostgreSQL 使用 MVCC。执行一次 UPDATE 时,数据库通常不是原地覆盖旧行,而是创建新版本,并把旧版本留给 VACUUM 回收。若一张表承担高频状态推进,例如:

pending -> processing -> retrying -> processing -> completed

同一条业务记录可能产生多个死亡元组。问题还会被以下因素放大:

  • 被更新的列出现在多个索引中,导致索引也需要持续写入新条目;
  • 长事务阻止 VACUUM 回收旧版本;
  • 大量工作进程竞争相同记录或页面;
  • 自动清理追不上写入速度;
  • 队列状态与业务实体混在同一张宽表中,每次更新都要改写较大的元组。

可以先用 PostgreSQL 自带统计视图筛查热点表:

SELECT
    schemaname,
    relname,
    n_live_tup,
    n_dead_tup,
    n_tup_ins,
    n_tup_upd,
    n_tup_del,
    n_tup_hot_upd,
    autovacuum_count,
    last_autovacuum
FROM pg_stat_user_tables
ORDER BY n_tup_upd DESC
LIMIT 20;

重点观察 n_tup_upd、n_dead_tup 和 autovacuum_count。如果更新次数极高、死亡元组持续增长,且自动清理非常频繁,就应该先检查数据模型和工作流,而不是立即提高 autovacuum 频率。

还要检查长事务,因为它们常常是“VACUUM 明明在运行,却回收不了空间”的原因:

SELECT
    pid,
    usename,
    application_name,
    state,
    now() - xact_start AS transaction_age,
    left(query, 120) AS query
FROM pg_stat_activity
WHERE xact_start IS NOT NULL
ORDER BY xact_start;

PgQ、PgQue 带来的设计转变

从架构角度看,这类工具的价值不只是“在 PostgreSQL 里放一个队列”,而是把隐含在业务表 UPDATE 语句里的流程显式化。

原来的系统可能通过轮询业务表工作:

SELECT *
FROM orders
WHERE status = 'pending'
ORDER BY created_at
LIMIT 100;

多个工作进程随后不断更新 orders.status、retry_count、worker_id 和时间戳。业务数据、调度状态和执行历史全部耦合在一起。

采用队列或工作流引擎后,可以改为:

  1. 业务事务写入订单,并产生一条待处理事件;
  2. 工作进程从队列领取事件,而不是扫描整张业务表;
  3. 重试策略、延迟执行和失败记录由队列层表达;
  4. 业务表只在真正的业务状态变化时更新;
  5. 历史过程通过事件记录保留,而不是反复覆盖同一列。

PgQ 与 PgQue 的具体安装方式、SQL API 和能力边界会随项目及版本变化,接入前应以目标版本文档为准。尤其要确认以下语义:

  • 消息是“至少一次”还是“至多一次”投递;
  • 消费者崩溃后,任务何时重新可见;
  • 是否支持延迟任务、批量消费和死信处理;
  • 消费者位置或租约如何持久化;
  • 是否能够与业务写入放在同一个数据库事务中。

最后一点尤其重要。若业务记录已经提交,但入队失败,就可能永久丢失任务;若先入队再写业务记录,消费者又可能读到尚不存在的数据。能够共享 PostgreSQL 事务,是数据库内队列的重要优势之一。

可以这样实践:构造一个最小任务队列

下面的示例不是 PgQ 或 PgQue 的正式 API,而是一个可直接在 PostgreSQL 中运行的最小模型,用来验证 SKIP LOCKED、幂等键和失败重试等设计。正式生产环境应优先评估成熟实现,不要把这个示例直接当成完整工作流引擎。

先创建表和索引:

CREATE TABLE workflow_job (
    id             bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    queue_name     text NOT NULL,
    idempotency_key text NOT NULL,
    payload        jsonb NOT NULL,
    available_at   timestamptz NOT NULL DEFAULT now(),
    attempts       integer NOT NULL DEFAULT 0,
    locked_at      timestamptz,
    locked_by      text,
    last_error     text,
    created_at     timestamptz NOT NULL DEFAULT now(),
    UNIQUE (queue_name, idempotency_key)
) WITH (fillfactor = 80);

CREATE INDEX workflow_job_ready_idx
    ON workflow_job (queue_name, available_at, id)
    WHERE locked_at IS NULL;

在业务事务中写数据并入队:

BEGIN;

INSERT INTO workflow_job (
    queue_name,
    idempotency_key,
    payload
)
VALUES (
    'email',
    'welcome-user-42',
    '{"user_id": 42, "template": "welcome"}'::jsonb
)
ON CONFLICT (queue_name, idempotency_key) DO NOTHING;

COMMIT;

工作进程可以在一个短事务中领取任务:

BEGIN;

WITH candidate AS (
    SELECT id
    FROM workflow_job
    WHERE queue_name = 'email'
      AND available_at <= now()
      AND locked_at IS NULL
    ORDER BY available_at, id
    FOR UPDATE SKIP LOCKED
    LIMIT 1
)
UPDATE workflow_job AS j
SET locked_at = now(),
    locked_by = 'worker-1',
    attempts = attempts + 1
FROM candidate
WHERE j.id = candidate.id
RETURNING j.id, j.payload, j.attempts;

COMMIT;

SKIP LOCKED 让多个消费者跳过已被其他事务锁住的任务,避免排队等待同一行。领取事务应尽量短:提交领取结果后再调用外部服务,不要在数据库事务中等待 HTTP 请求。

任务成功后删除队列记录:

DELETE FROM workflow_job
WHERE id = 123
  AND locked_by = 'worker-1';

失败时释放任务,并使用简单的退避策略:

UPDATE workflow_job
SET locked_at = NULL,
    locked_by = NULL,
    available_at = now() + make_interval(secs => LEAST(300, attempts * 10)),
    last_error = 'SMTP timeout'
WHERE id = 123
  AND locked_by = 'worker-1';

这套模型仍会产生 UPDATE 和 DELETE,因此仍然需要 VACUUM。它的意义是把写放大限制在窄而独立的队列表中,让核心业务表不再承担调度器职责。队列表也更容易单独设置 fillfactor、autovacuum 参数、分区和保留周期。

队列不是免维护区

把工作流放进 PostgreSQL 不等于问题自动消失。任务领取、重试和删除本身仍会制造死亡元组。如果流量很大,需要继续关注以下方面。

让索引保持克制

只为实际的领取查询建立索引。对 locked_by、attempts、last_error 等频繁变化的列随意加索引,会降低 HOT UPDATE 的机会并增加写放大。

缩短事务和外部调用的边界

推荐的消费者流程是:

短事务领取任务并提交
        ↓
事务外调用外部服务
        ↓
短事务确认成功或安排重试

这种模式通常意味着至少一次执行。消费者必须通过业务幂等键防止重复扣款、重复发邮件或重复创建资源。

为过期租约设计恢复机制

如果工作进程领取任务后崩溃,locked_at 不能永久阻止任务执行。可以定期释放超时租约:

UPDATE workflow_job
SET locked_at = NULL,
    locked_by = NULL,
    available_at = now()
WHERE locked_at < now() - interval '15 minutes';

超时时间必须大于正常任务处理时长,或者改用心跳续租。否则,一个仍在运行的任务可能被另一个消费者重复领取。

按表调整 autovacuum,而不是全局猛调

确认队列表确实需要更积极的清理后,可以只修改该表:

ALTER TABLE workflow_job SET (
    autovacuum_vacuum_scale_factor = 0.02,
    autovacuum_vacuum_threshold = 1000,
    autovacuum_analyze_scale_factor = 0.05
);

这些值只是起点,不是通用答案。应结合表大小、更新速率、磁盘吞吐和 VACUUM 实际耗时进行压测。若根因是长事务、无效索引或消费者设计错误,单纯调高 autovacuum 只会增加系统负担。

接入前的决策清单

PgQ、PgQue 或自建队列比较适合以下场景:

  • 任务与 PostgreSQL 中的业务数据需要原子提交;
  • 已有数据库运维能力,希望减少独立消息系统;
  • 队列吞吐量与保留周期处于 PostgreSQL 可承受范围;
  • 消费者能够实现幂等处理;
  • 团队愿意持续监控队列深度、任务年龄、重试次数和表膨胀。

如果系统需要跨地域超大规模事件流、长期保存全部消息、复杂流式处理或与大量异构系统解耦,专用消息平台通常更合适。

采用这类工作流引擎时,不要只比较“每秒能入队多少任务”。更有价值的指标包括:最老待处理任务的年龄、失败重试分布、重复执行率、死亡元组增长速度、VACUUM 耗时,以及消费者故障后的恢复时间。PgQ 和 PgQue 的真正价值,是让原本散落在 UPDATE、轮询脚本和定时任务中的流程变得可观察、可恢复,也更容易治理。


相关推荐