跳转至

消息队列:异步执行链与最终一致性

用户上传一份 PDF 文档,等待系统解析、向量化、建索引。如果这条链路全部塞进一次 HTTP 请求里同步执行,OCR 抖动或向量化模型超时都会让用户的浏览器转圈转到请求超时——而且失败之后没有任何补救路径,用户只能重新点一次上传,把整条链路的算力再消耗一遍。

消息队列(MQ)解决的正是这类"已经不适合继续占着这次请求,但必须保证最终会做完"的动作:把重型 I/O 从同步链路里挪出去,靠持久化日志保证执行的最终确定性。引入 MQ 后,这条文档处理链路获得了时序解耦与削峰填谷的能力,但也带来了分布式事务边界模糊(写库和发消息不是一个事务)、投递保证语义选择、消费幂等性设计等新的工程约束——接下来几节都会围绕这一条文档处理链路展开。

1. 异步边界:从同步链路到事件驱动

同步方案要求在单次 HTTP 请求内串行完成 保存记录 → OCR 解析 → 向量化 → 写索引,生产环境下任意一环的抖动都会造成请求超时,且失败后缺乏可靠的补偿路径。

  • 数据库定状态,MQ 推进度:请求进入系统后,首先在数据库以 TaskStatus: Pending 记录任务初始状态。状态写入成功后,系统即可向上游返回成功响应。后续的重型处理任务通过 MQ 投递至后台 Worker 异步执行。这意味着用户看到的"成功"是指"系统已承诺会完成",而非"已经完成"。
  • 副作用优先原则:消费端提交 Offset 的时序至关重要。正确顺序为:执行副作用 → 数据库更新状态 → 提交 Offset。若先提交 Offset 后执行业务逻辑,Worker 宕机将产生"消息已标记消费但数据未落库"的状态空洞——比如 OCR 结果丢了但 Offset 已经前移,这条消息不会被重新投递,数据永久丢失。

2. 系统容量模型与排队时延度量

消息系统的积压特征由生产消费速率差决定。当生产速率长期高于消费速率时,系统积压呈线性增长:

\[\text{Backlog Growth} \approx \text{Produce Rate} - \text{Consume Rate}\]

引入 Little's Law 变体,消息的平均排队等待时延可近似表示为:

\[\text{Queuing Delay} \approx \frac{\text{Backlog}}{\text{Consume Rate}}\]

这揭示了 MQ 的本质物理约束:消息队列不会缩短业务处理的总耗时,它将同步阻塞的线程占用转移为异步队列等待,从而释放上游系统的并发吞吐上限。 如果 OCR/向量化 Worker 的处理能力不足以消化上传速率(比如促销期间大量用户集中上传合同文件),积压只会不断增长,用户会发现"文档一直停在解析中"。

3. Kafka:基于分区日志的流式存储

Kafka 的本质是分布式分区追加日志系统,而非传统意义上的点对点消息队列。其顺序保证、并行度单位与消费语义均由此派生。

  • Partition 与局部顺序:顺序保证仅在单个 Partition 内部成立。要保证同一文档的解析任务按序执行(比如先完成 OCR 再触发向量化,不能反过来),必须以 document_id 作为 Partition Key,将同文档的所有消息路由至同一分区。全局顺序消费需强制约束在单分区,代价是放弃水平扩展能力。
  • Consumer Group 与并行度:Partition 是消费并行度的最小单位。Topic 的分区数决定了消费侧的有效并发上限——若文档解析 Topic 仅有 2 个分区,扩容至 10 个 Consumer 实例,实际参与消费的也只有 2 个,其余 8 个处于空闲状态。因此分区数是水平扩展的核心控制抓手
  • Rebalance(消费组重平衡):Consumer 加入、退出或心跳超时时,触发 Partition 的重新分配。此期间未提交 Offset 的消息会在新节点被重复拉取,是偶发重复消费的主要根因(比如同一份文档被重复 OCR 一次,浪费一次调用但不产生错误结果)。Rebalance 频率过高(如因 GC 停顿或网络抖动)会严重影响消费稳定性。

4. 投递保证语义(Delivery Guarantees)

语义 Offset 提交时机 消息丢失风险 重复消费风险 消费端要求
At-Most-Once 拉取后立即提交 有(Worker 宕机时丢失) 无特殊要求
At-Least-Once 业务逻辑完成并落盘后提交 有(Rebalance/重启) 必须实现幂等性
Effectively-Once At-Least-Once + 消费端强幂等去重 应用层感知无重复 幂等键 + 事务性写入

文档解析这类"宁可重复处理一次也不能丢"的链路只能选 At-Least-Once 或 Effectively-Once——At-Least-Once 是业界工程基石,Effectively-Once 在底层通常通过"At-Least-Once 投递 + 消费端数据库主键幂等去重"或"两阶段提交(2PC)与幂等写入事务数据库"来实现应用层感知的精准一次语义。

5. 消费治理:幂等、重试与死信

  • 幂等性是异步链路的入场条件:At-Least-Once 语义下,重复消费是常态而非异常。应通过数据库主键唯一约束或状态机条件更新(如 UPDATE ... WHERE status='Pending')封堵重复执行路径——同一份文档的 OCR 结果如果被重复消费两次却各写一份索引记录,用户搜索时会看到两条重复结果。任何不具备幂等保证的消费逻辑,在分布式环境下都是定时炸弹。
  • 重试预算分类治理
    • 临时故障(网络抖动、上游限流 429):进入重试队列,采用指数退避(Exponential Backoff)策略,如 1s → 2s → 4s → 8s,避免重试风暴(Retry Storm)压垮下游。
    • 永久故障(参数格式错误、权限不足、数据结构不兼容):重试无意义,应直接转入死信队列(DLQ),触发告警并等待人工介入处理——比如上传了一份加密 PDF 导致 OCR 无法解析,重试一万次结果都一样。DLQ 中的消息应包含完整的错误上下文和原始消息内容,以便定位根因。
  • Lag(消费积压)是核心健康指标:积压增长表明消费速率持续低于生产速率。初步响应通常为增加 Partition 数并同步扩容 Consumer 实例,同时排查单次消费耗时是否存在异常(如下游超时、锁等待)。
  • 幂等键贯穿始终:在实际工程中,位点更新、数据库状态修改、外部 API 调用、搜索引擎索引同步等多个外部副作用需保证最终一致性。应将幂等键(Idempotent Key)贯穿始终,结合本地消息表或分布式原子操作,防止状态重入带来的数据污染。

6. 消息系统与同步 RPC 选型对比

系统 核心优势 适用场景 局限
Kafka 高吞吐、追加日志回放、生态成熟 海量事件流、日志采集、RAG 异步建库(如文档解析这条链路) 单条消息延迟不是最优;运维复杂度较高
RocketMQ 事务消息、延时消息、业务治理能力强 复杂业务逻辑、金融级可靠性 社区生态相对集中于国内
RabbitMQ 细粒度路由、低延迟、协议灵活(AMQP) 任务分发、轻量级异步、RPC 场景 高吞吐场景性能受限;持久化开销较高
同步 HTTP/gRPC 结果立即可用,无中间件依赖 下游少、结果必须立即参与当前决策(如查询用户是否有权限上传) 下游故障会阻塞上游并产生级联失败

文档解析、权限变更、数据同步这类事件通常适合走消息队列,因为消费者可以独立扩展并从历史位点回放;查询类动作如果必须在当前请求内拿到结果,更适合同步 RPC。选择的核心依据是调用方是否必须在当前请求内获得结果。

7. 故障排查序列

  1. 查 Offset 与 Lag:进度滞留在哪个 Partition?是局部积压(某个文档 ID 热点)还是全局跟不上(Worker 处理能力不足)?
  2. 查副作用状态一致性:数据库中的 TaskStatus 与 MQ 消费进度是否对齐?是否存在大量文档的 TaskStatus 停留在 Processing 状态超过预期时间?
  3. 查 Rebalance 频率:Consumer 的 session.timeout.msmax.poll.interval.ms 配置是否合理?是否因 GC 停顿或单次消费耗时过长(比如某份文档特别大导致 OCR 超时)频繁触发重平衡?
  4. 查 Broker 资源:磁盘是否写满?PageCache 命中率是否下降导致读取性能退化?网络带宽是否被大消息占满?

核心结论Offset 是执行进度的标记,而非业务状态的真相。健壮的异步系统以数据库状态机作为最终收口,以 MQ 分区日志推进执行进度,并将幂等性作为消费端的基础设计约束,从而在重复投递的常态下保持正确性。

8. Outbox:封住"写库成功、发消息失败"的窗口

回到文档上传的例子:文档解析完成后,系统需要把"解析已完成"这件事同步通知给索引服务、审计日志和缓存失效队列。如果直接先写库再发 Kafka 消息,会在两个不能共同提交的系统之间留下窗口:

业务事务提交成功 → MQ 发送失败
MQ 发送成功       → 业务事务回滚

Transactional Outbox 把"文档已解析完成"和"需要通知下游"一起写进同一个本地事务。提交成功后,Relay 从 outbox 表可靠地把事件送往 Kafka:

sequenceDiagram participant API as DocumentService participant DB as MySQL participant R as Outbox Relay participant K as Kafka participant C as Consumer(索引/审计/缓存失效) API->>DB: 事务写解析结果 + outbox event DB-->>API: COMMIT R->>DB: 扫描未发送事件 R->>K: publish(event_id, version) K-->>R: ack R->>DB: 标记已发送 K->>C: 至少一次投递 C->>DB: 幂等写入 + 记录 event_id C->>K: commit offset

Relay 可能在 Kafka 已接收、数据库尚未标记发送完成时崩溃,重启后会再次投递——消费端的幂等处理原则见第 5 节,这里不再重复:Outbox 只负责不丢事件,重复由消费者兜底承受。

同一份文档的事件使用 document_id 做 Partition Key(呼应第 3 节的分区顺序约束),每条事件带 event_idaggregate_idversionoccurred_at。下游只接受比当前版本更新的事件,并写入回执;超时扫描和定期对账用于发现那些"事件已发送、效果未落地"的尾部失败——比如索引服务收到了事件但建索引失败,文档在搜索里始终不可见(这类"写库与写外部系统无法原子提交"的问题,在对象存储链路里也有对应形态,参见对象存储篇第 5.2 节)。

面试怎么答:面试官问同步改异步怎么落地,先说边界再说细节——用数据库状态机记录任务的权威状态,MQ 只负责推进执行进度,消费顺序必须是"先做事、再落库、最后提交 Offset",反过来会丢数据。问怎么保证消息不丢,说 Outbox:业务写库和写 outbox 表是同一个本地事务,Relay 异步转发,重复投递交给消费端幂等处理,因为 Outbox 解决的是"不丢",不是"不重复"。问 Kafka 和 RabbitMQ 怎么选,一句话说完:数据量大、要回放、要保序选 Kafka;路由复杂、要低延迟确认、任务分发选 RabbitMQ;调用方必须马上拿到结果就别用 MQ,直接同步 RPC。