You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka消费者集成测试故障求助:Testcontainers环境下消费失败

Kafka消费者Testcontainers集成测试修复方案

核心问题总结

你的测试存在三个关键问题:

  • Spring Boot消费者未配置连接到Testcontainers Kafka容器,仍使用默认的localhost:9092
  • 生产者与消费者的序列化/反序列化不匹配,发送的是字符串而非JSON对象
  • 容器启动时机和消息等待逻辑不够可靠,导致主题未就绪或消费未完成就断言

具体修复步骤

1. 让消费者连接到Testcontainers Kafka

在@DynamicPropertySource中添加Kafka bootstrap服务器配置,将Spring Boot消费者指向容器内的Kafka实例:

@DynamicPropertySource
public static void setProperties(DynamicPropertyRegistry registry) {
    registry.add("spring.datasource.url", mySQLContainer::getJdbcUrl);
    registry.add("spring.datasource.username", mySQLContainer::getUsername);
    registry.add("spring.datasource.password", mySQLContainer::getPassword);
    // 新增:动态注入Kafka容器的bootstrap地址
    registry.add("spring.kafka.bootstrap-servers", kafkaContainer::getBootstrapServers);
}

2. 修复生产者序列化逻辑

测试中的KafkaTemplate泛型错误,且发送了对象的字符串形式,导致消费者无法解析。修改如下:

// 修正泛型为<String, RecipientSavedEvent>
private static KafkaTemplate<String, RecipientSavedEvent> kafkaTemplate;

@BeforeAll
public static void setUp() {
    Map<String, Object> producerProps = new HashMap<>();
    producerProps.put("bootstrap.servers", kafkaContainer.getBootstrapServers());
    producerProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
    producerProps.put("value.serializer", "org.springframework.kafka.support.serializer.JsonSerializer");
    // 可选:配置类型映射,确保消费者识别对象类型
    producerProps.put(JsonSerializer.TYPE_MAPPINGS, 
        "com.tappedtechnologies.userservice.events.RecipientSavedEvent:com.tappedtechnologies.userservice.events.RecipientSavedEvent");

    kafkaTemplate = new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(producerProps));
    kafkaTemplate.setDefaultTopic("tappedtechnologies.emails.recipients");
}

// 测试方法中直接发送对象,无需toString()
@Test
public void consumePayload_Should_SavePayload() throws InterruptedException {
    RecipientSavedEvent expected = getPayload();
    kafkaTemplate.sendDefault(expected.getPayloadKey(), expected);
    // ... 后续断言逻辑
}

3. 清理冗余的消费者配置

application-test.properties中重复配置了反序列化器,删除无效行:

# Kafka Properties
spring.kafka.topic.name=tappedtechnologies.emails.recipients

# Kafka Consumer Properties
spring.kafka.consumer.group-id=userId
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties.spring.json.default.type=com.tappedtechnologies.userservice.events.RecipientSavedEvent
# 删掉下面这行冗余配置
# spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer

4. 移除手动容器启停代码

Testcontainers的@Container注解会自动管理容器生命周期,无需手动调用start()/stop():

@BeforeAll
public static void setUp() {
    // ... 仅保留KafkaTemplate初始化逻辑,删除以下两行
    // kafkaContainer.start();
    // mySQLContainer.start();
}

@AfterAll
public static void tearDown() {
    // 删除以下两行
    // kafkaContainer.stop();
    // mySQLContainer.stop();
}

5. 替换Thread.sleep为可靠等待

用CountDownLatch实现消费完成的精确等待(仅测试环境生效):
在消费者类中添加Latch:

@Slf4j
@Service
@RequiredArgsConstructor
@Profile("test") // 仅测试环境加载
public class RecipientPayloadConsumer implements KafkaConsumer<RecipientSavedEvent> {

    private final Creator<User> userCreator;
    public final CountDownLatch consumeLatch = new CountDownLatch(1);

    @Override
    @KafkaListener(topics = "${spring.kafka.topic.name}", groupId = "${spring.kafka.consumer.group-id}")
    public void consumePayload(@Payload RecipientSavedEvent payload) {
        log.info("Payload received from 'email-service': {}", payload);
        User userToSave = this.convertPayloadToUser(payload);
        try {
            userCreator.create(userToSave);
            consumeLatch.countDown(); // 消费完成后递减计数
        } catch (IllegalStateException exception) {
            log.error("IllegalStateException caught attempting to save payload from email-service.");
            log.error("Message: {}", exception.getMessage());
            throw exception;
        }
    }

    // ... 其他代码
}

测试方法中等待Latch:

@Test
public void consumePayload_Should_SavePayload() throws InterruptedException {
    RecipientSavedEvent expected = getPayload();
    kafkaTemplate.sendDefault(expected.getPayloadKey(), expected);
    
    // 等待消费完成,超时5秒
    assertThat(sut.consumeLatch.await(5, TimeUnit.SECONDS)).isTrue();

    User actual = userRepository.findByEmail(expected.getEmail()).orElse(null);

    assertThat(actual).isNotNull();
    assertThat(actual.getFirstName()).isEqualTo(expected.getFirstName());
    assertThat(actual.getLastName()).isEqualTo(expected.getLastName());
    assertThat(actual.getEmail()).isEqualTo(expected.getEmail());
}

6. 确保主题自动创建

添加配置让Kafka自动创建主题,避免LEADER_NOT_AVAILABLE警告:

# application-test.properties中新增
spring.kafka.admin.properties.auto.create.topics.enable=true

内容的提问来源于stack exchange,提问作者charbs29

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 00:39:52