使用 Kafka 作为消息队列时,有些坑是因为它设计初衷是"分布式日志系统"而非传统消息队列(如 RabbitMQ),如果照搬传统 MQ 的思维方式,很容易踩坑。以下是常见的几类问题:
1. 消息顺序问题
- 分区内有序,分区间无序:Kafka 只保证同一分区内的消息顺序,如果你的业务需要严格顺序(比如同一用户的操作按顺序处理),必须保证相关消息落到同一分区(合理设计 key)。
- 重试导致乱序:生产者开启重试(retries)且
max.in.flight.requests.per.connection > 1时,重试的消息可能"插队",导致同一分区内乱序。如果对顺序敏感,建议设置为 1,或使用enable.idempotence=true(幂等生产者会自动处理这个问题)。
2. 重复消费 / 消息丢失
- at-least-once 是默认语义:消费者处理完消息后才提交 offset,如果提交前挂掉,重启后会重复消费。业务逻辑必须做到幂等,否则会出问题(比如重复扣款)。
- 自动提交的陷阱:
enable.auto.commit=true时,offset 是按时间间隔自动提交的,可能在消息还没处理完就提交了(消息丢失),也可能处理完了还没提交就挂了(重复消费)。生产环境通常建议手动提交,并且在业务处理成功后再提交。 - acks 设置不当:
acks=0或acks=1在某些故障场景下会丢消息,需要金融级可靠性的场景建议acks=all配合min.insync.replicas。
3. 消费者组与 Rebalance 风暴
- Rebalance 影响吞吐:消费者数量变化、消费者处理超时(超过
max.poll.interval.ms)都会触发 rebalance,期间整个消费者组会短暂停止消费。如果单条消息处理耗时很长,容易频繁触发 rebalance。 - 消费者数量 > 分区数:多余的消费者会一直空闲,分区数决定了消费并行度的上限,扩容消费者之前要先看分区数够不够。
4. 把 Kafka 当成传统 MQ 用
- 没有真正的"消息确认/拒绝/重新入队"机制:不像 RabbitMQ 那样可以针对单条消息 ack/nack,Kafka 的消费进度是通过 offset 位移推进的,想做"这条失败了重试,其他继续"就得自己实现(比如死信队列、单独重试 topic)。
- 没有优先级队列、延迟队列的原生支持:这些都得自己在应用层实现(比如用时间戳 + 定时扫描,或额外引入延迟消息组件)。
5. 分区数设置一旦"想少了"很麻烦
- 分区数只能增加不能减少,而增加分区会破坏原有的 key-分区映射关系(导致同 key 消息路由到不同分区,破坏顺序性)。初期要预估好吞吐和并行度需求。
6. 消息堆积与磁盘问题
- 消费跟不上生产速度会导致 lag 持续增长,如果没有监控 lag 指标,问题往往是等到出大事才发现。
- 日志保留策略(retention):默认按时间或大小滚动删除,如果消费者故障时间过长,数据可能已经被删除,造成永久丢失(不像传统 MQ 消息消费后才删除)。
7. 生产端批量与延迟的权衡
linger.ms和batch.size设置不当,要么延迟高,要么吞吐低,需要根据业务场景(实时性 vs 吞吐量)调优。
8. 大消息问题
- Kafka 默认单条消息大小限制在 1MB 左右(
message.max.bytes),如果业务经常产生大消息(比如图片、大 JSON),需要提前调整配置或改用"消息体存对象存储、Kafka 只传引用"的模式。
总结一句话:Kafka 本质是个分布式提交日志(commit log),用它做消息队列时,很多"队列该有的功能"(单条确认、优先级、延迟)都得自己在应用层补上,同时顺序性、幂等性、offset 管理这几块是最容易踩坑、也最需要在设计初期就想清楚的地方。