1. Kafka java consumer为什么采用单线程的设计?

Kafka consumer其实是双线程,用户主线程和心跳线程。
处理消息的逻辑是在用户主线程完成的,从这个角度是单线程的。
Consumer获取到消息后,处理消息的逻辑是否采用多线程,由开发者决定。
所有的网络IO处理都发生在用户主线程,在多线程中共享同一个kafka consumer实例,会抛出异常。检查当前线程是否独占操作权。若检测到其他线程正在操作,则直接抛异常。

  1. 多线程方案
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)
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐