【问题标题】:Cannot disable manual commits on Kafka message using Spring-integration-kafka in Spring-Boot App无法在 Spring-Boot App 中使用 Spring-integration-kafka 禁用对 Kafka 消息的手动提交
【发布时间】:2017-01-06 02:02:55
【问题描述】:

由于我是 kafka 新手,我的团队正在 Spring-boot 应用程序上使用 Spring-Integration-kafka:2.0.0.RELEASE。基于此示例here,我能够使用我的 KafkaConsumer 使用 kafka 消息。

对于这个 spring-boot 应用程序,我有我的 Application.java

package hello;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import java.util.concurrent.TimeUnit;
import hello.notification.Listener;

@SpringBootApplication
public class Application {

    private Listener listener;

    public void run(String... args) throws Exception {
        this.listener.countDownLatch1.await(60, TimeUnit.SECONDS);
    }

    public static void main(String[] args) {
        SpringApplication.run(Application.class, args);
    }
}

这是我的 KafkaConsumerConfig.java

package hello.notification;

import ....; //different imports

@Configuration
@EnableKafka
public class KafkaConsumerConfig {
    @Bean
    KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        factory.setConcurrency(1);
        factory.getContainerProperties().setPollTimeout(3000);
        return factory;
    }

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> propsMap = new HashMap<>();
        propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, "group1");
        propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        //propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 100);
        propsMap.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 6000);
        propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        return new DefaultKafkaConsumerFactory<>(propsMap);
    }

    @Bean
    public Listener listener() {
        return new Listener();
    }
}

最后是我的 Listener.java

package hello.notification;

import ...; // different imports

public class Listener {

    public final CountDownLatch countDownLatch1 = new CountDownLatch(1);

    @KafkaListener(topics = "topic1")
    public void listen(ConsumerRecord<?, ?> record, Acknowledgment ack) {
        System.out.println(record + " and " + ack);
                // Place holder for "ack.acknowledge();"
        countDownLatch1.countDown();
    } 
}

我的问题是:

1) 在 ConsumerConfig 设置中,我已经将“ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG”设置为 false 并注释掉了“ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG”,为什么它仍然自动提交?(因为我不再会如果我重新启动我的 spring-boot 应用程序,请查看消息 - 我希望它会尝试重新发送,直到它被提交)

2) 为了手动提交,我想我只需要在我的 Listener.listen 中添加 ack.acknowledge() (请参阅占位符)?除了这个,我还需要其他什么来确保这是手动提交吗?

3)我想拥有最简单/最干净的 kafkaconsumer,我想知道你是否认为我所拥有的是最简单的方法,最初我正在寻找单线程消费者但那里没有很多示例所以我坚持并发监听器。

感谢您的所有帮助和投入!

【问题讨论】:

    标签: spring-boot apache-kafka spring-integration kafka-consumer-api spring-kafka


    【解决方案1】:

    为什么还是自动提交?

    因为默认情况下KafkaMessageListenerContainer 提供有:

    private AbstractMessageListenerContainer.AckMode ackMode = AckMode.BATCH;
    

    除此之外我还需要其他什么来确保这是手动提交吗?

    你必须切换到

    ContainerProperties.setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL)
    

    不知道你是否认为我所拥有的是最简单的方法

    没错。只需要习惯就好。

    更多信息请见Reference Manual

    【讨论】:

    • 感谢您的反馈。从您所说的来看,听起来我需要进行手动提交的唯一更改是在“KafkaConsumerConfig.java”文件中。我尝试添加“factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE);”在我返回工厂并重试但没有成功之前,重新启动我的消息仍然没有重新出现。请问您是否可以使用我上面的代码提供一些示例?谢谢!
    • 好吧,尝试将propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); 添加到consumerFactory 道具。默认为latestkafka.apache.org/documentation/#consumerconfigs
    • 我添加了 propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");并且在 Listener 类中,我们确保我们实现了 AcknowledgeingMessageListener 并覆盖了 onMessage 方法,但仍然没有让消息返回的运气,感觉就像它在 onMessage() 结束时自动提交:公共类 Listener 实现 AcknowledgeingMessageListener{ @Override public void onMessage(ConsumerRecord, ?> record, Acknowledgment ack) { System.out.println(record + " and " + ack); } }
    • 我假设您的意思是您希望在下一次轮询迭代中再次收到该消息,但 Kafka 并非如此。消费者中有一些内部状态来保存日志的index。应用重启后即可再次获取。
    • 是的,有点..我的目标是创建一个干净的 KafkaConsumer,除非我告诉它(基于条件并使用 ack.acknowledge()),否则它不会提交。所以我创建了一个类似于上面的代码,为了测试消费者必须由我手动提交,我故意不使用 ack.acknowledge() 来强制它在下次重启时再次向我重新发送相同的消息,但我还没有找到解决方案。虽然我没有提交它(我没有使用 ack.acknowledge()),但当我重新启动我的应用程序时,味精没有回来......这意味着它已经被提交......这有意义吗?任何样本都可以提供?谢谢
    猜你喜欢
    • 2018-07-17
    • 2018-11-23
    • 1970-01-01
    • 2019-08-15
    • 2016-09-04
    • 1970-01-01
    • 2020-07-10
    • 2018-10-03
    • 2021-03-21
    相关资源
    最近更新 更多