【发布时间】:2022-01-05 12:38:23
【问题描述】:
我有 2 个 Spring Boot 应用程序,一个是 Kafka 发布者,另一个是消费者。我正在尝试编写集成测试以确保发送和接收事件。
在 IDE 中或从命令行运行时测试为绿色,而无需其他测试,例如 mvn test -Dtest=KafkaPublisherTest。但是,当我构建整个项目时,测试以org.awaitility.core.ConditionTimeoutException 失败。项目中有多个@EmbeddedKafka 测试。
在日志中的这些行之后测试卡住了:
2021-11-30 09:17:12.366 INFO 1437 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : wages-consumer-test: partitions assigned: [wages-test-0, wages-test-1]
2021-11-30 09:17:14.464 INFO 1437 --- [er-event-thread] kafka.controller.KafkaController : [Controller id=0] Processing automatic preferred replica leader election
如果您对如何测试这些东西有更好的想法,请分享。
这是测试的样子:
@SpringBootTest(properties = { "kafka.wages-topic.bootstrap-address=${spring.embedded.kafka.brokers}" })
@EmbeddedKafka(partitions = 1, topics = "${kafka.wages-topic.name}")
class KafkaPublisherTest {
@Autowired
private TestWageProcessor testWageProcessor;
@Autowired
private KafkaPublisher kafkaPublisher;
@Autowired
private KafkaTemplate<String, WageEvent> kafkaTemplate;
@Test
void publish() {
Date date = new Date();
WageCreateDto wageCreateDto = new WageCreateDto().setName("test").setSurname("test").setWage(BigDecimal.ONE).setEventTime(date);
kafkaPublisher.publish(wageCreateDto);
kafkaTemplate.flush();
WageEvent expected = new WageEvent().setName("test").setSurname("test").setWage(BigDecimal.ONE).setEventTimeMillis(date.toInstant().toEpochMilli());
await()
.atLeast(Duration.ONE_HUNDRED_MILLISECONDS)
.atMost(Duration.TEN_SECONDS)
.with()
.pollInterval(Duration.ONE_HUNDRED_MILLISECONDS)
.until(testWageProcessor::getLastReceivedWageEvent, equalTo(expected));
}
}
发布者配置:
@Configuration
@EnableConfigurationProperties(WagesTopicPublisherProperties.class)
public class KafkaConfiguration {
@Bean
public KafkaAdmin kafkaAdmin(WagesTopicPublisherProperties wagesTopicPublisherProperties) {
Map<String, Object> configs = new HashMap<>();
configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, wagesTopicPublisherProperties.getBootstrapAddress());
return new KafkaAdmin(configs);
}
@Bean
public NewTopic wagesTopic(WagesTopicPublisherProperties wagesTopicPublisherProperties) {
return new NewTopic(wagesTopicPublisherProperties.getName(), wagesTopicPublisherProperties.getPartitions(), wagesTopicPublisherProperties.getReplicationFactor());
}
@Primary
@Bean
public WageEventSerde wageEventSerde() {
return new WageEventSerde();
}
@Bean
public ProducerFactory<String, WageEvent> producerFactory(WagesTopicPublisherProperties wagesTopicPublisherProperties, WageEventSerde wageEventSerde) {
Map<String, Object> configProps = new HashMap<>();
configProps.put(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
wagesTopicPublisherProperties.getBootstrapAddress());
configProps.put(
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class);
configProps.put(
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
wageEventSerde.serializer().getClass());
return new DefaultKafkaProducerFactory<>(configProps);
}
@Bean
public KafkaTemplate<String, WageEvent> kafkaTemplate(ProducerFactory<String, WageEvent> producerFactory) {
return new KafkaTemplate<>(producerFactory);
}
}
消费者配置:
@Configuration
@EnableConfigurationProperties(WagesTopicConsumerProperties.class)
public class ConsumerConfiguration {
@ConditionalOnMissingBean(WageEventSerde.class)
@Bean
public WageEventSerde wageEventSerde() {
return new WageEventSerde();
}
@Bean
public ConsumerFactory<String, WageEvent> wageConsumerFactory(WagesTopicConsumerProperties wagesTopicConsumerProperties, WageEventSerde wageEventSerde) {
Map<String, Object> props = new HashMap<>();
props.put(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
wagesTopicConsumerProperties.getBootstrapAddress());
props.put(
ConsumerConfig.GROUP_ID_CONFIG,
wagesTopicConsumerProperties.getGroupId());
props.put(
ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class);
props.put(
ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
wageEventSerde.deserializer().getClass());
return new DefaultKafkaConsumerFactory<>(
props,
new StringDeserializer(),
wageEventSerde.deserializer());
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, WageEvent> wageEventConcurrentKafkaListenerContainerFactory(ConsumerFactory<String, WageEvent> wageConsumerFactory) {
ConcurrentKafkaListenerContainerFactory<String, WageEvent> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(wageConsumerFactory);
return factory;
}
}
@KafkaListener(
topics = "${kafka.wages-topic.name}",
containerFactory = "wageEventConcurrentKafkaListenerContainerFactory")
public void consumeWage(WageEvent wageEvent) {
log.info("Wage event received: " + wageEvent);
wageProcessor.process(wageEvent);
}
这里是项目源码:https://github.com/aleksei17/springboot-rest-kafka-mysql
以下是构建失败的日志:https://drive.google.com/file/d/1uE2w8rmJhJy35s4UJXf4_ON3hs9JR6Au/view?usp=sharing
【问题讨论】:
-
看起来您正在启动多个引导应用程序和多个代理;考虑对所有测试使用单个代理实例 - docs.spring.io/spring-kafka/docs/current/reference/html/…
-
@GaryRussell,我不知道这样的选项,谢谢!但它不适用于我的情况,这里是代码:github.com/aleksei17/springboot-rest-kafka-mysql/commit/… 所以我会检查使用 Spring Cloud Stream 是否会使代码更具可测试性
标签: spring-boot apache-kafka spring-kafka spring-kafka-test