【问题标题】:spring-kafka: Kafka consumer isn't receiving records if I define a record producer in my applicationspring-kafka:如果我在我的应用程序中定义记录生产者,Kafka 消费者不会收到记录
【发布时间】:2017-12-01 21:46:21
【问题描述】:

我正在使用 Spring 和 Spring Kafka 编写一个小型 PoC。我的目标是同时拥有生产者和消费者,从这个主题写入(分别读取)。

我遇到了一个奇怪的情况:

  • 生产者正在正确生成记录(我可以通过 Python 脚本使用它们)
  • 消费者没有收到记录
  • 但是如果我从我的代码中删除生产者并通过另一种方法(例如使用 python 脚本)生成记录,消费者会正确接收记录。

以下是我的代码 - 它与文档中的示例非常相似。更准确地说,问题出在 KafkaConsumerConfiguration 中的 bean 不是由 Spring 创建的事实(即从未调用构造它们的方法)。

制作人

KafkaProducerConfiguration.java

@Configuration
public class KafkaProducerConfiguration {

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:32768");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return new DefaultKafkaProducerFactory<>(props);
    }
}

MessageSender.java

@Component
public class MessageSender {
    final static private Logger log = Logger.getLogger(MessageSender.class);

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @PostConstruct
    public void onConstruct() throws InterruptedException {
        log.info("Sending messages...");
        for (int i = 0; i < 100; ++i) {
            kafkaTemplate.send("mytopic", "this is a message");
            Thread.sleep(1000);
        }
        kafkaTemplate.flush(); // NOTE: no changes if I move this call in the loop
        log.info("Done sending messages");
    }
}

消费者

KafkaConsumerConfiguration.java

@Configuration
@EnableKafka
public class KafkaConsumerConfiguration {

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

        return factory;
    }

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

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:32768");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-service");
        return props;
    }
}

MyMessageListener.java

@Service
public class MyMessageListener {
    final static private Logger log = Logger.getLogger(MyMessageListener.class);

    @PostConstruct
    public void onConstruct() {
        log.info("Message listener started");
    }

    @KafkaListener(topics = "mytopic")
    public void onMessageReceived(String message) {
        log.info("Got message: "+ message);
    }
}

这是应用程序生成的日志,供参考:https://pastebin.com/BY783jiL。如您所见,消费者 bean 没有被创建(否则会出现一个块 ConsumerConfig values: ...

以下是我尝试过的一些事情,但没有成功:

  • 将生产者和消费者配置放在同一个配置类中
  • 更改KafkaConsumerConfiguration中bean的名称(并在MyMessageListener.onMessageReceived方法上添加注解属性containerFactory = "myBeanName"
  • 将类的名称 KafkaConsumerConfiguration 更改为其他名称
  • 在我的 KafkaConsumerConfiguration 中添加一个与 kafka 无关的 @Bean 以查看它是否会被创建(确实如此)

版本:Spring Boot 1.5.9,Spring-Kafka:1.1.7。

我已经把头发扯了几个小时了,感谢任何帮助。

谢谢!

【问题讨论】:

    标签: spring spring-boot apache-kafka spring-kafka


    【解决方案1】:
    kafkaTemplate.send("mytopic", "this is a message");
    

    您永远不应该在 @PostConstruct 方法中开始与外部服务交互 - 您需要等待应用程序构建完成后再这样做。

    实现SmartLifecyle,为isAutoStartup 返回true,并将该代码移至start()

    或者实现 ApplicationListener&lt;ConstextRefreshedEvent&gt; 并在收到事件时进行发送。

    任何一种方式都将确保应用程序准备就绪。

    【讨论】:

      【解决方案2】:

      刚刚发现问题。 MessageSender.onConstruct 实际上需要很长时间来执行(100 秒),同时它会阻止 Spring 创建其他 bean。

      【讨论】:

        猜你喜欢
        • 2019-11-21
        • 2021-04-25
        • 1970-01-01
        • 2018-06-20
        • 1970-01-01
        • 2019-05-09
        • 2023-04-08
        • 2017-02-25
        • 1970-01-01
        相关资源
        最近更新 更多