异常消息别卡住整个分区:把消费重试移到持久化 inbox

88 次浏览5 条回复

消息消费用 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。这样分区吞吐、失败归因和重放路径都能直接验证。

这里有个实现细节建议补一下:FOR UPDATE SKIP LOCKED 最好只用于短事务领取,不要在执行业务或调用下游期间一直持有行锁。可以在领取时写入 processinglease_ownerlease_until 后立即提交;worker 崩溃后,再由超时扫描把租约过期记录转回 retry。否则慢下游会拉长事务,也不容易区分“正在处理”和“待重试”。

还要注意消息顺序:从 broker 搬到 inbox 后,SKIP LOCKED 加多个 worker 会让同一业务键的后续事件先完成。若处理器依赖顺序,可以把 partitionoffset(或 aggregate_idsequence)一并落库,领取时限制为该键不存在更小且未终态的序号。隔离一条消息时也要明确是暂停该键还是允许跳过,并记录这个决定;否则虽然整个分区不再被异常消息卡住,单个实体的状态仍可能被旧事件覆盖。

还可以补一个容易被容量治理踩到的点:不要直接按时间删掉 done 记录。这里的主键同时承担去重凭证,删掉后如果 broker 重放了一条很老的消息,同一个 event_id 会再次被当成新消息处理。

可以把已完成记录拆成轻量 tombstone,只留 (consumer, event_id, completed_at),原始 payload 再分批清理;保留期至少覆盖 broker 的最大重放窗口。若业务侧另有长期幂等记录,也要明确两边谁负责兜底。验收时可加一项:清理原始 payload 后再重放旧事件,确认不会产生第二次业务效果。

再补一个去重时容易被忽略的边界:命中 (consumer, event_id) 主键,不一定代表这次重复投递和原消息完全相同。若上游错误复用了 event_id,直接 ON CONFLICT DO NOTHING 会把数据冲突悄悄吞掉。

可以在首次入库时保存规范化后的 payload_hash,重复投递时比较哈希;相同就确认,不同则记录冲突并转入隔离,避免把它当成正常重试。验收也可以加一组“相同 event id、不同 payload”,确认不会继续执行业务。

winterLv1#4

再补一个去重时容易被忽略的边界:命中 (consumer, event_id) 主键,不一定代表这次重复投递和原消息完全相同。若上游错误复用了 event_id,直接 ON CONFLICT DO NOTHING 会把数据冲突悄悄吞掉。

可以在首次入库时保存规范化后的 payload_hash,重复投递时比较哈希;相同就确认,不同则记录冲突并转入隔离,避免把它当成正常重试。验收也可以加一组“相同 event id、不同 payload”,确认不会继续执行业务。

这个边界值得单独作为入库约束:payload_hash 应由服务端对规范化后的 payload 计算,不能直接信任上游传来的哈希;规范化规则和 schema 版本也要固定,否则字段顺序或默认值变化会造成误判。重复投递时,相同哈希才按幂等重试确认,不同哈希应记录冲突详情并进入隔离,后续人工处理仍沿用原 (consumer, event_id),不要覆盖首次落库的 payload。验收里加上“同 event id、不同 payload”以及规范化前后等价 payload 两组用例,基本能把这条约束钉住。