【问题标题】:Spring Kafka : Record listener vs Batch listenerSpring Kafka:记录侦听器与批处理侦听器
【发布时间】:2022-05-03 21:51:03
【问题描述】:

使用 spring-kafka,有两种类型的 Kafka 监听器。

Record Listeners

@KafkaListener(groupId = "group1", topics = {"my.topic"})
public void listenSingle(String message, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) {
    /* Process my kafka message */
}

还有Batch Listeners

/*
    Consumer factory is initialized with setBatchListener(true)
*/

@KafkaListener(groupId = "group1", topics = {"my.topic"})
public void listenBatch(List<String> messages, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) throws Exception {
    messages.forEach({
        /* Process my kafka message */
    });
}

根据文档,它似乎对 Kafka 消费者没有任何影响(无论如何都会轮询多条消息)。

然后我不明白为什么要使用批处理侦听器而不是其他侦听器,因为批处理侦听器有一些记录侦听器没有的限制(拦截器、偏移管理等)?

也许我误解了什么?批处理侦听器有什么好处?

【问题讨论】:

  • 也许是因为您可以确认整个批次而不是单个消息?
  • @cricket_007 说的很好,逐个确认消息对性能的影响真的很大吗?
  • 可以,但要看你是至少要一次,还是最多一次

标签: java spring apache-kafka spring-kafka


【解决方案1】:

在我的用例中,优势不是来自对 Kafka 的处理,而是来自侦听器中的后续处理。例如,如果您必须在侦听器的消息处理中调用 REST API,您可以使用批处理侦听器以批量方式执行此操作。您可以将整个列表传递给 API。当然,外部 API 也必须支持批量操作。另一个示例可能是针对在处理 Kafka 记录时访问的数据库进行批量处理。

例如:

@KafkaListener
public void receive(List<AnyPojo> pojo) {
  myPojoRepository.saveAll(pojo);
}

如果您在没有批处理的情况下执行此操作,这将导致每条记录都有一个新事务,这比批量/批量执行要慢得多:

@KafkaListener
public void receive(AnyPojo pojo) {
  myPojoRepository.save(pojo);
}

【讨论】:

    猜你喜欢
    • 2021-08-31
    • 1970-01-01
    • 2018-03-07
    • 2019-08-25
    • 2022-11-11
    • 2011-08-16
    • 2021-02-03
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多