消息的顺序怎么保证?消息积压了怎么处理?
一句话回答
全局有序要求单分区、单消费者,吞吐上不去,实际都做局部有序:只保证同一个业务键(比如同一个订单)的消息有序。做法是发送时按业务键选择分区或队列,让同一个键的消息进同一个分区;消费时同一个键串行处理,可以单线程消费一个分区,也可以在消费者内按键分组并行。积压时先找原因再处理:消费者故障或变慢就修复、扩容(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 的新消息可能进入另一个分区,而它的旧消息还留在原分区。在原分区的积压消费完之前,新消息可能先被处理,所以要求严格有序的主题扩分区前要先停写或等积压清空。
代码示例
消费者内部按业务键分组:同一个键串行,不同键并行。
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 轮
这道题你掌握了吗?
选一个最接近的状态,没掌握的题会出现在"我的进度 · 待复习"里。
学习记录暂存在本机浏览器。登录后自动同步到账号,换设备也能看到。