【问题标题】:How to use Spring Kafka's Acknowledgement.acknowledge() method for manual commit如何使用 Spring Kafka 的 Acknowledgement.acknowledge() 方法进行手动提交
【发布时间】:2021-11-09 22:53:07
【问题描述】:

我第一次使用 Spring Kafka,我无法在我的消费者代码中使用 Acknowledgement.acknowledge() 方法进行手动提交,如此处https://docs.spring.io/spring-kafka/reference/html/_reference.html#committing-offsets 所述。我的是弹簧启动应用程序。如果我不使用手动提交过程,那么我的代码就可以正常工作。但是当我使用 Acknowledgement.acknowledge() 用于手动提交,它显示与 bean 相关的错误。另外,如果我没有正确使用手动提交,请建议我正确的方法。

错误信息:

***************************
APPLICATION FAILED TO START
***************************

Description:

Field ack in Receiver required a bean of type 'org.springframework.kafka.support.Acknowledgment' that could not be found.


Action:

Consider defining a bean of type 'org.springframework.kafka.support.Acknowledgment' in your configuration.

我在谷歌上搜索了这个错误,发现我需要添加@Component,但这已经在我的消费者代码中了。

我的消费者代码如下所示:Receiver.java

import java.util.concurrent.CountDownLatch;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;

@Component
public class Receiver {

    @Autowired
    public Acknowledgment ack;

    private CountDownLatch latch = new CountDownLatch(1);

    @KafkaListener(topics = "${kafka.topic.TestTopic}")
    public void receive(ConsumerRecord<?, ?> consumerRecord){
            System.out.println(consumerRecord.value());
            latch.countDown();
            ack.acknowledge();
    }
}

我的生产者代码如下所示:Sender.java

import java.util.Map;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;

@Component
public class Sender {

    @Autowired
    private KafkaTemplate<String, Map<String, Object>> kafkaTemplate;

    public void send(Map<String, Object> map){
            kafkaTemplate.send("TestTopic", map);

    }

}

编辑 1:

我的新消费者代码如下所示:Receiver.java

import java.util.concurrent.CountDownLatch;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;

@Component
public class Receiver {

    private CountDownLatch latch = new CountDownLatch(1);

    @KafkaListener(topics = "${kafka.topic.TestTopic}", containerFactory = "kafkaManualAckListenerContainerFactory")
    public void receive(ConsumerRecord<?, ?> consumerRecord, Acknowledgment ack){
            System.out.println(consumerRecord.value());
            latch.countDown();
            ack.acknowledge();
    }
}

我也改变了我的配置类:

import java.util.HashMap;
import java.util.Map;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

@Configuration
@EnableKafka
public class ReceiverConfig {

    @Value("${kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Value("${spring.kafka.consumer.group-id}")
    private String consumerGroupId;

    @Bean
    public Map<String, Object> consumerConfigs() throws SendGridException {
            Map<String, Object> props = new HashMap<>();
            props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
            props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
            return props;

    }

    @Bean
    public ConsumerFactory<String, String> consumerFactory(){
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    /*@Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(){
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());

        return factory;
    }*/

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaManualAckListenerContainerFactory(){
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

    @Bean
    public Receiver receiver() {
        return new Receiver();
    }
}

将 containerFactory = "kafkaManualAckListenerContainerFactory" 添加到我的 receive() 方法后,我收到以下错误。

***************************
APPLICATION FAILED TO START
***************************

Description:

Parameter 1 of method kafkaListenerContainerFactory in org.springframework.boot.autoconfigure.kafka.KafkaAnnotationDrivenConfiguration required a bean of type 'org.springframework.kafka.core.ConsumerFactory' that could not be found.
    - Bean method 'kafkaConsumerFactory' in 'KafkaAutoConfiguration' not loaded because @ConditionalOnMissingBean (types: org.springframework.kafka.core.ConsumerFactory; SearchStrategy: all) found bean 'consumerFactory'


Action:

Consider revisiting the conditions above or defining a bean of type 'org.springframework.kafka.core.ConsumerFactory' in your configuration.

【问题讨论】:

    标签: java spring-boot kafka-consumer-api spring-kafka


    【解决方案1】:

    对于那些仍在寻找解决这些与手动确认有关的错误的方法的人,您无需指定 containerFactory = "kafkaManualAckListenerContainerFactory",而只需添加:

    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
    

    在您返回工厂对象之前到您的接收器配置。

    那你还需要:

    props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
    

    在消费者配置道具中。

    所以最后你的监听器方法可以简单地看起来像:

    @KafkaListener(topics = "${spring.kafka.topic}")
        private void listen(@Payload String payload, Acknowledgment acknowledgment) {
            //Whatever code you want to do with the payload
            acknowledgement.acknowledge(); //or even pass the acknowledgment to a different method and acknowledge even later
        }
    

    【讨论】:

    • 有没有人有一个使用这种模式的整个工作应用程序的最小工作示例?
    • 当然,我为你创建了这个简单的项目,看看我的 git 页面:github.com/zim8662/kafka-example 只需在本地运行 zookeeper 和 kafka 服务器,或者如果它是远程更改 application.properties 中的属性(当然还有组和主题属性)并在浏览器中尝试 localhost:8080/test。
    • 嗨 Zim - 想知道如果你设置 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true) 会发生什么? kafka 还会自动提交吗? Spring 会再次尝试提交吗?
    • @alex 我认为这是对这个属性的一个简单但很好的解释:medium.com/@danieljameskay/…
    • 我认为您的回答不好,因为更改将影响绑定到 ListenerContainerFactory 的所有 Kafka 侦听器。
    【解决方案2】:

    你真的应该关注documentation

    使用手动AckMode时,也可以给监听器提供Acknowledgment;这个例子还展示了如何使用不同的容器工厂。

    @KafkaListener(id = "baz", topics = "myTopic",
              containerFactory = "kafkaManualAckListenerContainerFactory")
    public void listen(String data, Acknowledgment ack) {
        ...
        ack.acknowledge();
    }
    

    确实没有注意到Acknowledgment 是一个bean。因此,请适当地更改您的 receive() @KafkaListener 方法签名并删除该 @Autowired 用于可疑的 Acknowledgment bean - 它只是不存在,因为此对象是每个接收到的消息的一部分(标头)。

    【讨论】:

    • 嗨@Artem,我做了更改,但现在又遇到了另一个 bean 错误。我用新的编辑更新我的问题。 spring kafka 文档不是很清楚。
    • 1.您的 receive() 方法中仍然没有 Acknowledgment arg。所以,我不确定ack var 的用途是什么。 2. 你有一些原始的 Spring Kafka 配置和 Spring Boot 的组合。您应该考虑从那里开始:docs.spring.io/spring-boot/docs/1.5.7.RELEASE/reference/… 并通过正确的 ConsumerFactory bean 重新配置开箱即用的 kafkaListenerContainerFactory bean。在您的情况下,它是 ConsumerFactory&lt;String, String&gt;,但必须是 ConsumerFactory&lt;?, ?&gt;
    • 说“你应该遵循文档”是徒劳的;这里的文档非常混乱,对初学者不友好。一个指向最小工作应用程序的链接会很好。
    • @ArtemBilan 文档链接已损坏。
    【解决方案3】:

    那些正在使用spring boot应用程序的人,只需将下面的内容添加到您的application.yml(或环境特定文件)中即可。

    spring:
      kafka:
        listener:
          ack-mode: manual
    

    以上更改将使Acknowledgment ack 参数在内部可用
    receive(ConsumerRecord&lt;?, ?&gt; consumerRecord, Acknowledgment ack) 方法自动。

    【讨论】:

      猜你喜欢
      • 2011-07-31
      • 2018-11-23
      • 1970-01-01
      • 2017-09-10
      • 2018-05-05
      • 2019-08-27
      • 1970-01-01
      • 1970-01-01
      • 2016-09-04
      相关资源
      最近更新 更多