Kafka共享消费者不消费消息
我尝试测试新的Share Consumer功能,在阅读了Spring Kafka的文档后,创建了一个 示例项目,但我无法按预期让它工作。
我的示例基于Spring Boot 4.1.0-RC1,并将 kafka.version 设置为 4.2.0 以使用最新的Kafka客户端,并且还使用 kafka:latest 的Docker镜像来确保testcontainers中 Kafka的版本为4.2。
配置类:
@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));
}
}
在运行测试时,消息已经成功发送,然而未被监听器接收。
The GreetingListener 是一个简单的Kafka监听器。
@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);
}
}
解决方案
因此,我最初使用 这个项目 在Docker中本地运行Kafka,并尝试共享分组。
然后,经过反复试验,我最终将你的测试成功通过:
- 我在Testcontainer中添加了
KAFKA_SHARE_COORDINATOR_STATE_TOPIC_REPLICATION_FACTOR属性 - 我在测试一开始就加入
Thread.sleep(5_000),因为根据日志,似乎在消息发送时消费者并未处于活跃状态,因此我决定在发布消息之前再等一会儿
另外,供参考,我注意到如果把配置中的 NewTopic myTopic Bean完全移除,测试并不会受到影响——它仍然通过(显然主题会自动创建)。
我理解 Thread.sleep 看起来像是一个权宜之计,可能有更优雅的解决方案,我也不能给出发生了什么或为什么会这样的百分之百清晰解释。不过,我希望我的回答至少提供了一个可行的解决方案和一个有用的方向,帮助你继续前进。
最终结果:
import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.kafka.KafkaContainer;
import org.testcontainers.utility.DockerImageName;
import java.time.Duration;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
@SpringBootTest
@Testcontainers
@Slf4j
class DemoApplicationTests {
@Container
static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("apache/kafka:4.2.0"))
.withEnv("KAFKA_SHARE_COORDINATOR_STATE_TOPIC_REPLICATION_FACTOR", "1"); // without this the test will fail
@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
void testSendMessage() throws Exception {
Thread.sleep(5_000); // wait for all kafka stuff is started
List.of("the", "quick", "brown", "fox", "jumps", "over", "the", "lazy", "dog")
.forEach(word -> kafkaTemplate.send(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));
}
}
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。