【发布时间】:2017-07-30 22:03:12
【问题描述】:
我正在尝试配置 Spring Boot 应用程序以使用 Kafka 消息。添加后:
<!-- https://mvnrepository.com/artifact/org.springframework.kafka/spring-kafka -->
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<version>1.1.3.RELEASE</version>
</dependency>
进入我的依赖项并使用@EnableKafka 和@KafkaListener(topics = "some-topic") 注释,我收到以下错误:
...
Caused by: org.springframework.beans.factory.NoSuchBeanDefinitionException: No bean named 'kafkaListenerContainerFactory' available
然后我添加以下配置:
@Bean
public Map<String, Object> consumerConfigs() {
Map<String, Object> propsMap = new HashMap<>();
propsMap.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
propsMap.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
propsMap.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "100");
propsMap.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
propsMap.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
propsMap.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
propsMap.put(ConsumerConfig.GROUP_ID_CONFIG, "group1");
propsMap.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
return propsMap;
}
@Bean
public ConsumerFactory<String, String> consumerFactory() {
return new DefaultKafkaConsumerFactory<>(consumerConfigs());
}
@Bean
KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.setConcurrency(3);
factory.getContainerProperties().setPollTimeout(3000);
return factory;
}
错误消失了。但是,我认为我应该能够使用 spring.kafka.listener.* 属性自动配置它,正如文档所建议的那样。
如果我不能,我想使用自动连接的KafkaProperties。但是,为了能够使用它,我添加:
<!-- https://mvnrepository.com/artifact/org.springframework.boot/spring-boot-autoconfigure -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-autoconfigure</artifactId>
<version>1.5.2.RELEASE</version>
</dependency>
然后就可以导入了。当我尝试如下使用它时:
@Autowired
private KafkaProperties kafkaProperties;
在我的方法中:
return kafkaProperties.buildConsumerProperties();
我收到以下错误:
Caused by: java.lang.ClassNotFoundException: org.springframework.boot.context.annotation.DeterminableImports
.
我认为这是一个 Maven 依赖问题。
所以我的问题是:
- 是否可以在不创建
@Beans 而仅使用application.properties的情况下配置 Kafka 配置? - 如果没有,我如何跳过手动创建所需的
Map对象,而直接使用kafkaProperties.buildConsumerProperties()而不会出现上述错误(第二个)?
【问题讨论】:
标签: spring spring-boot apache-kafka spring-integration