Kafka使用RoundRobin分区器时分区分布不均问题咨询
两次调用
partition方法的设计原因 Kafka Producer 2.8版本doSend逻辑中的两次partition调用是官方有意设计的容错机制:
- 第一次调用是为了初步计算消息归属分区,校验当前分区是否存在可用的未满批次可以直接写入该消息
- 当判定需要新建批次时,Producer会先主动拉取最新的主题元数据(避免第一次调用时用了过期元数据,分配到已经失效、或者leader不可用的分区),再第二次调用
partition方法拿到基于最新元数据计算的分区号,再创建对应批次,本质是为了提升分区分配的准确性和可用性。
RoundRobin分区器不均匀的根因
官方自带的2.8版本RoundRobin分区器是有状态的,内部维护了一个自增计数器,每调用一次partition方法计数器就会自增1,两次调用就会直接跳过一个分区编号,长期运行就会出现固定分区永远分配不到消息的问题,无法达到均匀分布的预期。
你自定义分区器方案的可行性
你当前的自定义逻辑是生产可用的,这也是业界针对该问题的通用修复思路,和Kafka 3.0及以上版本官方对RoundRobin分区器的修复逻辑一致:新增幂等判断,同一个消息触发的多次partition调用不会重复触发计数器自增,直接返回第一次计算的分区号即可。
上线前只要注意两个细节即可:
- 分区器内部的计数器必须使用
AtomicInteger等线程安全的原子类实现,避免多线程发送场景下的计数错乱 - 仅对同一个消息的重复调用返回上次的分区号,不要跨不同消息复用分区编号,避免破坏轮询逻辑
生产环境可用的均匀分区方案
- 方案一:升级Kafka客户端版本到3.0及以上,该版本官方已经直接修复了RoundRobin分区器的重复调用跳号问题,无需修改业务代码,兼容性和稳定性最高
- 方案二:如果无法升级客户端,直接复用3.0版本官方的
RoundRobinPartitioner实现代码,自定义为本地分区器引入即可,官方实现已经做了完整的线程安全、幂等处理,完全可以在2.8版本的服务端集群正常使用 - 方案三:如果不需要严格的逐消息轮询,也可以直接使用Kafka默认的Sticky分区器,该分区器会尽量将同一批次的消息写入同一个分区减少请求开销,批次写满后自动切换下一个分区,长期运行也能实现分区的均匀分布,性能比普通RoundRobin更高
内容的提问来源于stack exchange,提问作者Indraneel Bende
相关产品推荐
相关产品推荐

