【发布时间】:2020-11-11 22:04:19
【问题描述】:
由于原生 KafkaConsumer 不是线程安全的,因此不鼓励从不同线程而不是 kafka-consumer 处理线程调用 pause 和 resume 方法。 但由于 spring-kafka 提供了另一层 KafkaMessageListenerContainer,它在内部使用 kafka-consumer。所以我的问题是我们是否可以使用 KafkaListenerEndpointRegistry 通过 id 获取侦听器容器并从其他线程而不是消费者处理线程调用恢复或暂停方法。
kafkaListenerEndpointRegistry.getListenerContainer("id").pause();
ExecutorService executorService = newFixedThreadPool(2);
executorService.submit(()->{
System.out.println("CurrentThread: {}" + Thread.currentThread().getId()+ " " + Thread.currentThread().getName());
try {
Thread.sleep(5000);
} catch (InterruptedException e) {
e.printStackTrace();
}
kafkaListenerEndpointRegistry.getListenerContainer("id").resume();
});
【问题讨论】:
标签: java apache-kafka spring-kafka