kafka消费者之多线程方案
·
- Kafka java consumer为什么采用单线程的设计?
Kafka consumer其实是双线程,用户主线程和心跳线程。
处理消息的逻辑是在用户主线程完成的,从这个角度是单线程的。
Consumer获取到消息后,处理消息的逻辑是否采用多线程,由开发者决定。
所有的网络IO处理都发生在用户主线程,在多线程中共享同一个kafka consumer实例,会抛出异常。检查当前线程是否独占操作权。若检测到其他线程正在操作,则直接抛异常。
- 多线程方案
kafkaTemplate.send(topic, orderId % partitions, message); // 按订单 ID 分区[1,7](@ref)
方案一:多线程+多kafka consumer实例
在一个消费者组中,每个分区都只能被组内的一个消费者实例所消费。假设一个消费者组订阅了100个分区,那么他只能扩展到100个线程。
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> factory() {
ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConcurrency(3); // 设置并发消费者数(需小于等于分区数)[3,6,8](@ref)
return factory;
}
方案二:单线程+单kakfa consumer实例+消息处理worker线程池
高伸缩性,如果消息获取速度慢,增加获取消息的线程数;如果消息处理速度慢,增加消息处理的线程。
最大的缺陷是获取消息和处理消息分开了,不是同一个线程处理了,因此无法保证顺序。而且线程池来消费会加长消费链路,导致位移提交变得困难,可能导致重复消费。
有序性和重复性都无法保证。
public void listen(ConsumerRecords<String, String> records) {
records.forEach(record ->
executor.execute(() -> process(record))); // 异步处理[1,3,7](@ref)
}
更多推荐




所有评论(0)