【发布时间】: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