Spring Kafka Share Consumer测试异常:消息发送成功但未被接收
基于Spring Boot 4.1.0-RC1搭建项目测试Spring Kafka新的Share Consumer功能,指定kafka.version为4.2.0,使用Testcontainers的apache/kafka:latest镜像部署Kafka 4.2环境。测试时消息发送成功,但GreetingListener无法接收消息。相关代码如下:
配置类
@Configuration @Slf4j class ShareConsumerConfig { @Value("${spring.kafka.bootstrap-servers}") String bootstrapServers; @Bean NewTopic myTopic() { return new NewTopic(DEMO_TOPIC_NAME, 1, (short) 1); } @Bean public ShareConsumerFactory<String, String> shareConsumerFactory() { log.debug("Get bootstrap servers from properties:{}", bootstrapServers); Map<String, Object> props = Map.of( ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers, ConsumerConfig.GROUP_ID_CONFIG, DEMO_GROUP_NAME, ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class, ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class //ConsumerConfig.SHARE_ACKNOWLEDGEMENT_MODE_CONFIG, "explicit" ); DefaultShareConsumerFactory<String, String> factory = new DefaultShareConsumerFactory<>(props); factory.addListener(new ShareConsumerFactory.Listener<>() { @Override public void consumerAdded(String id, ShareConsumer<String, String> consumer) { log.debug("consumer added id:{}", id); } @Override public void consumerRemoved(@Nullable String id, ShareConsumer<String, String> consumer) { log.debug("consumer removed id:{}", id); } }); return factory; } @Bean public ShareKafkaListenerContainerFactory<String, String> shareKafkaListenerContainerFactory( ShareConsumerFactory<String, String> shareConsumerFactory) { return new ShareKafkaListenerContainerFactory<>(shareConsumerFactory); } }
测试类
@Testcontainers @SpringBootTest @Slf4j class DemoApplicationTests { // Kafka 4.2 enabled share consumer by default @Container static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("apache/kafka:latest")); @DynamicPropertySource static void kafkaProperties(DynamicPropertyRegistry registry) { registry.add("spring.kafka.bootstrap-servers", kafkaContainer::getBootstrapServers); } @Autowired private KafkaTemplate<String, String> kafkaTemplate; @Autowired private GreetingListener listener; @Test public void testSendMessage() { List.of("the", "quick", "brown", "fox", "jumps", "over", "the", "lazy", "dog") .forEach(word -> kafkaTemplate.send(DemoApplication.DEMO_TOPIC_NAME, word) .thenAccept(s -> log.debug("sent message: {}", s))); Awaitility.waitAtMost(Duration.ofMillis(30_000)) .untilAsserted(() -> assertThat(this.listener.getWordCount("the")).isEqualTo(2)); } }
监听器类
@Component @Slf4j public class GreetingListener { public Map<String, Long> counter = new ConcurrentHashMap<>(); @KafkaListener( topics = DEMO_TOPIC_NAME, containerFactory = "shareKafkaListenerContainerFactory", groupId = DEMO_GROUP_NAME ) public void onMessage(ConsumerRecord<String, String> record) { log.debug("received record: {} at {}", record, LocalDateTime.now()); counter.compute(record.value(), (s, v) -> v == null ? 1 : v + 1); } public Long getWordCount(String word) { return this.counter.get(word); } }
内容的提问来源于stack exchange,提问作者Hantsy
相关产品推荐
相关产品推荐

