【发布时间】:2022-06-10 22:25:18
【问题描述】:
我的 Spring Boot 应用程序将每小时从 kafka 代理收听 100 万条记录。每条消息的整个处理逻辑需要 1-1.5 秒,包括数据库插入。 Broker有64个分区,这也是我@KafkaListener的并发。
在我每小时收听大约 50k 条记录的较低环境中,我当前的代码只能在一分钟内处理 90 条记录。以下是代码,所有其他配置参数,如 max.poll.records 等都是默认值:
@KafkaListener(id="xyz-listener", concurrency="64", topics="my-topic")
public void listener(String record) {
// processing logic
}
我确实得到每小时 7-8 次“消费者很可能被踢出群组”。我认为这两个问题都可以通过隔离侦听器方法和对每条消息进行多线程处理来解决,但我不知道该怎么做。
【问题讨论】:
标签: spring multithreading spring-kafka