【问题标题】:Kafka Consumer not getting invoked when the kafka Producer is set to Sync当卡夫卡生产者设置为同步时,卡夫卡消费者没有被调用
【发布时间】:2017-02-22 01:45:48
【问题描述】:

我有一个要求,其中需要维护 2 个主题,其中 1 个采用同步方法,另一种采用异步方法。 异步调用消费者记录按预期工作,但是在同步方法中,消费者代码没有被调用。

下面是配置文件中声明的代码

 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9093");
 props.put(ProducerConfig.RETRIES_CONFIG, 3);
 props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
 props.put(ProducerConfig.ACKS_CONFIG, "all");
 props.put(ProducerConfig.LINGER_MS_CONFIG, 1);
 props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);

我在这里启用了 autoFlush true

 @Bean( name="KafkaPayloadSyncTemplate")
    public KafkaTemplate<String, KafkaPayload> KafkaPayloadSyncTemplate() {
        return new KafkaTemplate<String,KafkaPayload>(producerFactory(),true);
 }

控件在返回recordMetadataResults对象后停止,不再对消费者进行任何调用

  private List<RecordMetadata> sendPayloadToKafkaTopicInSync() throws   InterruptedException, ExecutionException {      
        final List<RecordMetadata> recordMetadataResults = new ArrayList<RecordMetadata>();
        KafkaPayload kafkaPayload = constructKafkaPayload();
        ListenableFuture<SendResult<String,KafkaPayload>> 
future = KafkaPayloadSyncTemplate.send(TestTopic, kafkaPayload);
        SendResult<String, KafkaPayload> results;
        results = future.get();
        recordMetadataResults.add(results.getRecordMetadata());     
        return recordMetadataResults;           
    }

消费者代码

public class KafkaTestListener {    
    @Autowired
    TestServiceImpl TestServiceImpl;    
    public final CountDownLatch countDownLatch = new CountDownLatch(1); 
    @KafkaListener(id="POC", topics = "TestTopic", group = "TestGroup")
    public void listen(ConsumerRecord<String,KafkaPayload> record, Acknowledgment acknowledgment) {
        countDownLatch.countDown();     
        TestServiceImpl.consumeKafkaMessage(record);        
        System.out.println("Acknowledgment : " + acknowledgment);
        acknowledgment.acknowledge();       
    }
}

基于问题,我有 2 个问题

  1. 当监听器类是同步生产者时,我们是否应该手动调用监听器类中的listen()。如果是,该怎么做?
  2. 如果侦听器 (@KafkaListener) 被自动调用,我需要添加哪些其他设置/配置才能使其正常工作。

提前感谢您的输入

-斯里坎特

【问题讨论】:

    标签: apache-kafka spring-kafka


    【解决方案1】:

    您应该确保将consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); 用于消费者属性。

    不确定您所说的同步/异步是什么意思,但生产和消费是完全不同的操作。而且您不能从生产者方面影响消费者。因为介于两者之间的是 Kafka Broker。

    【讨论】:

    • 这个属性已经添加到 consumerprops 中了。就我而言,消息没有被消耗。你的意思是,无论我是否声明我的生产者类型在同步模式或异步模式下工作,消费者操作都应该工作?
    • 同步模式:我期待一个 ack,因此我声明了 KafkaTemplate 与 autoFlush 设置为真。异步模式:我调用future的回调方法。
    • 当然对面话题有没有消费者也没关系。生产者只是向主题发送消息。在这种情况下,同步/异步意味着您如何等待确认消息存储在主题中。这里绝对没有关于消费者的信息。不知道发生了什么。也许你可以分享一些简单的 Spring Boot 应用程序,我们会看看它是否存在缺陷。
    • 我不知道最初出了什么问题,但它现在工作,可能是其他一些依赖问题。感谢您的回复,它现在正在工作。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-08-06
    • 1970-01-01
    • 1970-01-01
    • 2019-07-03
    • 2018-05-05
    相关资源
    最近更新 更多