【发布时间】:2019-06-24 14:44:27
【问题描述】:
我在 Spring Boot 中有 kafka 处理程序:
@KafkaListener(topics = "topic-one", groupId = "response")
public void listen(String response) {
myService.processResponse(response);
}
例如生产者每秒发送一条消息。但是myService.processResponse 工作 10 秒。我需要处理每条消息并在新线程中启动myService.processResponse。我可以创建我的执行者并将每个响应委托给它。但我认为 kafka 中还有其他配置供他们使用。我找到了 2:
1) 将 concurrency = "5" 添加到 @KafkaListener 注释 - 它似乎正在工作。但我不确定有多正确,因为我有第二种方法:
2) 我可以创建ConcurrentKafkaListenerContainerFactory 并将其设置为ConsumerFactory 和concurrency
我不明白这些方法之间的区别?只需将concurrency = "5" 添加到@KafkaListener 注释就足够了,还是我需要创建ConcurrentKafkaListenerContainerFactory?
或者我什么都不懂,还有其他办法吗?
【问题讨论】:
标签: java spring-boot apache-kafka kafka-consumer-api spring-kafka