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

使用Kafka与Testcontainers做集成测试遇认证问题求助

问题:Kafka容器集成测试中生产者认证失败报错

我正在为应用开发集成测试,其中一个测试流程使用Kafka向客户端通知更新。已创建KafkaContainer实现类:

import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.containers.KafkaContainer;
import org.testcontainers.utility.DockerImageName;

import java.time.Duration;

public class MyKafkaContainer extends KafkaContainer {

    private static final String KAFKA_DOCKER_IMAGE_NAME = "my_image_name";

    private static KafkaContainer kafkaContainer;

    private MyKafkaContainer() {
        super(DockerImageName.parse(KAFKA_DOCKER_IMAGE_NAME));
    }

    public static KafkaContainer getInstance() {
        if (kafkaContainer == null) {
            kafkaContainer = new KafkaContainer(DockerImageName.parse(KAFKA_DOCKER_IMAGE_NAME).asCompatibleSubstituteFor("confluentinc/cp-kafka"))
                    .withStartupAttempts(3)
                    .withStartupTimeout(Duration.ofMinutes(3));
        }
        return kafkaContainer;
    }

    @Override
    public void start() {
        super.start();
    }

    @Override
    public void stop() {
        //do nothing, JVM handles shut down
    }
}

随后创建初始化类设置bootstrap-servers参数:

public class TestContainersInitializer implements ApplicationContextInitializer<ConfigurableApplicationContext> {


    private static final KafkaContainer kafkaContainer;


    static {

        kafkaContainer = MyKafkaContainer.getInstance();


        Startables.deepStart(kafkaContainer).join();
    }

    @Override
    public void initialize(ConfigurableApplicationContext applicationContext) {
        TestPropertyValues.of(
                "spring.kafka.bootstrap-servers=" + kafkaContainer.getBootstrapServers(),
        ).applyTo(applicationContext.getEnvironment());
    }
}

在测试类添加注解:@ContextConfiguration(initializers = TestContainersInitializer.class)

启动正常,但发送Kafka消息时,实例化生产者出现错误:

org.apache.kafka.common.KafkaException: javax.security.auth.login.LoginException: Could not login: the client is being asked for a password, but the Kafka client code does not currently support obtaining a password from the user. not available to garner authentication information from the user.


原因分析

这个错误的核心是:你的Kafka客户端尝试使用SASL认证连接测试容器,但测试用的Kafka容器默认未开启认证,而你的应用继承了生产环境的认证配置(比如security.protocol=SASL_PLAINTEXT),导致客户端发起认证请求却没有配置有效凭证,最终抛出登录异常。

测试环境不需要像生产环境那样提供密码,Testcontainers启动的Kafka容器默认是无认证的PLAINTEXT模式。

解决方法

1. 显式禁用测试环境的Kafka认证

在测试配置中覆盖认证相关属性,强制使用PLAINTEXT协议:

  • 方式一:在测试类添加@TestPropertySource注解
@TestPropertySource(properties = {
    "spring.kafka.properties.security.protocol=PLAINTEXT",
    "spring.kafka.properties.sasl.jaas.config=" // 清空JAAS认证配置
})
  • 方式二:在TestContainersInitializer中添加属性配置
@Override
public void initialize(ConfigurableApplicationContext applicationContext) {
    TestPropertyValues.of(
            "spring.kafka.bootstrap-servers=" + kafkaContainer.getBootstrapServers(),
            "spring.kafka.properties.security.protocol=PLAINTEXT",
            "spring.kafka.properties.sasl.mechanism=",
            "spring.kafka.properties.sasl.jaas.config="
    ).applyTo(applicationContext.getEnvironment());
}

2. 检查自定义Kafka镜像的默认配置

你使用的是自定义镜像my_image_name,如果该镜像默认开启了SASL认证,需要在启动容器时添加环境变量关闭认证:

public static KafkaContainer getInstance() {
    if (kafkaContainer == null) {
        kafkaContainer = new KafkaContainer(DockerImageName.parse(KAFKA_DOCKER_IMAGE_NAME).asCompatibleSubstituteFor("confluentinc/cp-kafka"))
                .withStartupAttempts(3)
                .withStartupTimeout(Duration.ofMinutes(3))
                .withEnv("KAFKA_SECURITY_PROTOCOL", "PLAINTEXT")
                .withEnv("KAFKA_SASL_ENABLED_MECHANISMS", "")
                .withEnv("KAFKA_AUTHORIZER_CLASS_NAME", "");
    }
    return kafkaContainer;
}

3. 修复单例容器的实例不一致问题

你的MyKafkaContainer单例实现存在逻辑问题:构造方法继承了自定义镜像,但getInstance却直接创建了KafkaContainer实例,而非MyKafkaContainer实例,可能导致配置不统一。建议修改为:

public static KafkaContainer getInstance() {
    if (kafkaContainer == null) {
        kafkaContainer = new MyKafkaContainer()
                .withStartupAttempts(3)
                .withStartupTimeout(Duration.ofMinutes(3));
    }
    return kafkaContainer;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 18:24:54