消息消费用 broker 原地重试,遇到一条永久无法处理的数据时,后面的正常消息也会一直排队。直接确认并丢弃又不利于修复后重放。一个比较小的做法是:broker 只负责把消息可靠写进本地 inbox,业务重试由独立 worker 完成。
create table consumer_inbox (
consumer text not null,
event_id text not null,
payload jsonb not null,
state text not null check (state in ('pending', 'retry', 'done', 'quarantined')),
attempts integer not null default 0,
next_attempt_at timestamptz,
error_code text,
handler_version text not null,
updated_at timestamptz not null default now(),
primary key (consumer, event_id)
);
收到消息时只做一次短事务:按 (consumer, event_id) 插入 inbox,提交成功后再确认 broker。重复投递命中主键即可确认。worker 用 for update skip locked 领取到期任务;临时故障按总时长和次数预算退避,格式或业务校验这类确定性错误直接进 quarantined,不要继续占重试容量。
如果业务写入和 inbox 在同一个数据库里,可以把业务变更与 done 放进同一事务。跨系统调用仍可能出现“对方成功、本地没标记”的窗口,所以调用时要用稳定的幂等键,例如 consumer:event_id,不能因为第几次尝试而变化。人工重放也建议显式把原记录从 quarantined 改回 pending 并记录新的 handler_version,不要悄悄换一个 event id 绕过去重。
验收可以连续投递正常、永久异常、正常、临时失败四条消息,再在业务提交前后分别杀掉 worker。后面的正常消息不应被异常消息挡住;重启后不能产生额外业务效果;临时失败应在预算内恢复或进入隔离;修复处理器后,隔离记录可以按原 event id 重放并最终到 done。这样分区吞吐、失败归因和重放路径都能直接验证。