消息的顺序怎么保证?消息积压了怎么处理?

进阶高频场景题约 7 分钟读完

一句话回答

全局有序要求单分区、单消费者,吞吐上不去,实际都做局部有序:只保证同一个业务键(比如同一个订单)的消息有序。做法是发送时按业务键选择分区或队列,让同一个键的消息进同一个分区;消费时同一个键串行处理,可以单线程消费一个分区,也可以在消费者内按键分组并行。积压时先找原因再处理:消费者故障或变慢就修复、扩容(Kafka 的消费者数受分区数限制);流量突增来不及扩分区,可以临时把消息转发到分区更多的新主题,用更多消费者处理;同时降级非核心逻辑、改成批量处理,并对 消费延迟(lag) 配置告警。

详细解析

为什么会乱序

文本
生产端:消息 1、2、3 发往不同分区 ──> 各分区独立消费,进度不同
消费端:一个分区的消息交给线程池并行处理 ──> 后到的先处理完
重试:消息 1 处理失败稍后重试,消息 2 先处理成功
生产端重试:开启多个在途请求又没开幂等时,第 1 批失败重试、第 2 批先写入

Kafka 只保证单个分区内有序;RabbitMQ 单个队列内投递有序,但多个消费者竞争消费、消息被 nack 后重新入队,都会打乱处理顺序。

局部有序怎么做

发送端:同一个键进同一个分区

  • Kafka:发送时指定 key,默认分区器对 key 做哈希后对分区数取模,同一个 key 总是落在同一个分区
  • RocketMQ:发送顺序消息时,用 MessageQueueSelector 按业务键选择队列
  • RabbitMQ:同一个键路由到同一个队列,这个队列只由一个消费者处理

生产者开启 enable.idempotence=true,Kafka 在重试时也能保证单分区内不乱序。

消费端:同一个键串行处理

  • 最简单:一个分区一个线程,按顺序处理,吞吐受限于分区数
  • 想提高并发:消费者拉到消息后按 key 哈希分到 N 个内存队列,每个队列一个线程串行处理。这时提交 offset 要小心,只能提交到"该分区所有更早的消息都已处理完"的位置
  • RocketMQ 的顺序消费模式会锁定队列,同一时刻只有一个线程消费一个队列

处理失败:顺序消费时,一条消息失败要原地重试,不能跳过它先处理后面的消息,否则就乱序了;重试多次仍失败,只能告警并暂停这个键(或整个分区),或者把它和同一个键的后续消息一起转入死信,人工处理。

另一种思路:不依赖消息顺序,消费端用版本号或状态机兜底,丢弃比当前状态更旧的消息,见 幂等。很多场景这样做比强行保证顺序更简单。

消息积压怎么处理

1. 先定位原因,看消费 lag 曲线和消费者日志:

现象 可能原因
lag 持续增长,消费速度掉到 0 消费者挂了、卡死,或者频繁 rebalance
消费速度下降 下游数据库或接口变慢,消费逻辑里有慢查询
消费速度正常,生产速度突然变大 活动流量、上游 bug 导致重复发送
某个分区 lag 特别高 热点 key 集中到一个分区,或者这个分区的消费者有问题

2. 按原因处理:

  • 消费者故障:修复后重启;如果是某条消息一直处理失败卡住了队列,把它转入死信,先放行后面的消息
  • 处理慢:优化慢查询、批量写库、提高单个消费者内的并发(注意上面讲的顺序和 offset 问题)
  • 扩容消费者:Kafka 中一个分区同时只能被组内一个消费者消费,消费者数超过分区数没有意义
  • 分区不够又急需扩容:新建一个分区数更多的临时主题,写一个只转发、不处理的程序把积压的消息搬过去,再部署大量消费者消费临时主题;处理完后切回原来的架构
  • 降级:暂停非核心的处理逻辑,或者对时效性已经过去的消息(比如过期的通知)直接丢弃,事后再补

3. 事后预防:对 lag 配置告警(绝对值和增长速度),提前按峰值规划分区数,压测出单个消费者的处理能力。

Kafka 可以给主题增加分区,但不能减少;增加分区后 key 到分区的映射会变,同一个 key 的新消息可能进入另一个分区,而它的旧消息还留在原分区。在原分区的积压消费完之前,新消息可能先被处理,所以要求严格有序的主题扩分区前要先停写或等积压清空。

代码示例

消费者内部按业务键分组:同一个键串行,不同键并行。

JavaScript
class KeyedExecutor {
  constructor(lanes = 8) {
    this.tails = Array.from({ length: lanes }, () => Promise.resolve()) // 每条通道的队尾
  }

  lane(key) {
    let h = 0
    for (const ch of String(key)) h = (h * 31 + ch.charCodeAt(0)) >>> 0
    return h % this.tails.length
  }

  submit(key, task) {
    const i = this.lane(key)
    const run = this.tails[i].then(task) // 排在同一通道上一个任务之后
    this.tails[i] = run.catch(() => {}) // 前一个失败不影响后面继续排队
    return run
  }
}

// 同一个订单的消息按到达顺序处理,不同订单并行
const executor = new KeyedExecutor(4)
for (const msg of batch) {
  executor.submit(msg.orderId, () => handle(msg))
}

面试官可能追问

一个 key 的消息特别多,导致单个分区积压怎么办?

这是热点 key 问题。可以把热点 key 拆开,比如在 key 后面加上子维度(订单 ID 加商品 ID),只要业务只要求更细粒度的有序就行;如果业务确实要求这个 key 全部有序,单分区的处理能力就是上限,只能优化单条消息的处理速度。

积压的消息过期被删掉了怎么办?

Kafka 按保留时间删除日志,RabbitMQ 消息可能因 TTL 进入死信或被丢弃。消息已经没了,只能从源头补:从上游的数据库或日志里查出这段时间的数据,重新生成消息发送。所以积压告警的阈值要远小于保留时间。

消费者频繁 rebalance 导致积压怎么排查?

常见原因是处理一批消息的时间超过了 max.poll.interval.ms,或者心跳超时(session.timeout.ms)。可以减小 max.poll.records、优化处理速度;容器频繁重启也会触发 rebalance。Kafka 2.4 起支持增量协作式 rebalance(把 partition.assignment.strategy 设为 CooperativeStickyAssignor),只迁移需要变动的分区,其他分区不停止消费;Kafka 4.0 起新的消费者组协议(group.protocol=consumer)正式可用,由服务端计算分配,进一步减少全组停顿。静态成员(group.instance.id)也能避免容器重启引起的 rebalance。

易错点

  • 只在生产端按 key 分区,消费端却用线程池乱序处理,最后还是乱序
  • 为了顺序把所有消息发到一个分区,吞吐上不去
  • 积压时只想着加消费者,没注意 Kafka 消费者数超过分区数就没用了
  • 顺序消费中跳过失败的消息继续处理,破坏了顺序

AI 模拟面试官

用自己的话回答,AI 对照参考答案打分、指出遗漏,再追问,最多 3 轮

登录后就可以和 AI 面试官对练,面试记录也会保存下来。登录

这道题你掌握了吗?

选一个最接近的状态,没掌握的题会出现在"我的进度 · 待复习"里。

学习记录暂存在本机浏览器。登录后自动同步到账号,换设备也能看到。