生产环境里,有一类表看起来不大,却能把 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 和时间戳。业务数据、调度状态和执行历史全部耦合在一起。
采用队列或工作流引擎后,可以改为:
- 业务事务写入订单,并产生一条待处理事件;
- 工作进程从队列领取事件,而不是扫描整张业务表;
- 重试策略、延迟执行和失败记录由队列层表达;
- 业务表只在真正的业务状态变化时更新;
- 历史过程通过事件记录保留,而不是反复覆盖同一列。
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、轮询脚本和定时任务中的流程变得可观察、可恢复,也更容易治理。