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

SpringBoot SSL连接Docker部署Kafka报API_VERSIONS请求取消问题

SpringBoot集成Kafka SSL连接Docker部署Broker失败问题

我搭建了一个集成简易Kafka生产者、消费者的SpringBoot应用,尝试通过SSL协议连接Docker部署的Kafka Broker,但始终无法正常运行,相关配置如下:

相关配置

docker-compose.yml

version: '3.1'

services:

  zookeeper:
    image: wurstmeister/zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181

  broker:
    image: wurstmeister/kafka
    ports:
      - 9092:9092
      - 9093:9093
    volumes:
      - ./security:/etc/kafka/secrets
    environment:
      KAFKA_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092,TLS://broker:29093,TLS_HOST://localhost:9093
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092,TLS://broker:29093,TLS_HOST://localhost:9093
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT,TLS:SSL,TLS_HOST:SSL
      KAFKA_INTER_BROKER_LISTENER_NAME: TLS
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'false'
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_SSL_KEYSTORE_LOCATION: /etc/kafka/secrets/broker-keypair.pem
      KAFKA_SSL_KEY_PASSWORD: password
      KAFKA_SSL_KEYSTORE_TYPE: PEM
      KAFKA_SSL_TRUSTSTORE_LOCATION: /etc/kafka/secrets/root.crt
      KAFKA_SSL_TRUSTSTORE_TYPE: PEM
      KAFKA_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM: " "
      KAFKA_SSL_CLIENT_AUTH: "required"
      KAFKA_ALLOW_EVERYONE_IF_NO_ACL_FOUND: 'true'

KafkaConsumerConfig.java

@EnableKafka
@Configuration
public class KafkaConsumerConfig {

    @Value(value = "${spring.kafka.properties.bootstrap.servers}")
    private String bootstrapAddress;

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        final Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "sales.quotation.quotation.receiver");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);

        props.put(
                SslConfigs.SSL_KEYSTORE_LOCATION_CONFIG,
                "path\\to\\consumer.keystore.jks");
        props.put(SslConfigs.SSL_KEYSTORE_PASSWORD_CONFIG, "password");
        props.put(SslConfigs.SSL_KEY_PASSWORD_CONFIG, "password");
        props.put(
                SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG,
                "path\\to\\truststore.jks");
        props.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "password");

        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        final ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

KafkaProducerConfig.java

@Configuration
public class KafkaProducerConfig {

    @Value(value = "${spring.kafka.properties.bootstrap.servers}")
    private String bootstrapAddress;

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        final Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);

        configProps.put(
                SslConfigs.SSL_KEYSTORE_LOCATION_CONFIG,
                "path\\to\\producer.keystore.jks");
        configProps.put(SslConfigs.SSL_KEYSTORE_PASSWORD_CONFIG, "password");
        configProps.put(SslConfigs.SSL_KEY_PASSWORD_CONFIG, "password");
        configProps.put(
                SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG,
                "path\\to\\truststore.jks");
        configProps.put(SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "password");

        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

application.properties

spring.kafka.properties.bootstrap.servers=${KAFKA_BOOTSTRAP_SERVERS}

测试类中通过如下代码设置配置项:

System.setProperty("KAFKA_BOOTSTRAP_SERVERS", "localhost:9093");

报错信息

所有证书均由根CA签发,Docker容器可正常启动,但运行SpringBoot测试类尝试连接Broker时抛出如下错误:

16:16:48.786 [kafka-admin-client-thread | adminclient-1] INFO  o.apache.kafka.clients.NetworkClient:341 - [AdminClient clientId=adminclient-1] Cancelled in-flight API_VERSIONS request with correlation id 26 due to node -1 being disconnected (elapsed time since creation: 2ms, elapsed time since send: 2ms, request timeout: 3600000ms)
16:16:48.893 [kafka-admin-client-thread | adminclient-1] INFO  o.apache.kafka.clients.NetworkClient:935 - [AdminClient clientId=adminclient-1] Node -1 disconnected.

16:17:42.434 [org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1] INFO  o.apache.kafka.clients.NetworkClient:935 - [Consumer clientId=consumer-sales.quotation.quotation.receiver-1, groupId=sales.quotation.quotation.receiver] Node -1 disconnected.
16:17:42.434 [kafka-producer-network-thread | producer-1] INFO  o.apache.kafka.clients.NetworkClient:341 - [Producer clientId=producer-1] Cancelled in-flight API_VERSIONS request with correlation id 150 due to node -1 being disconnected (elapsed time since creation: 1ms, elapsed time since send: 1ms, request timeout: 30000ms)
16:17:42.434 [org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1] INFO  o.apache.kafka.clients.NetworkClient:341 - [Consumer clientId=consumer-sales.quotation.quotation.receiver-1, groupId=sales.quotation.quotation.receiver] Cancelled in-flight API_VERSIONS request with correlation id 153 due to node -1 being disconnected (elapsed time since creation: 2ms, elapsed time since send: 2ms, request timeout: 30000ms)
16:17:42.434 [kafka-producer-network-thread | producer-1] WARN  o.apache.kafka.clients.NetworkClient:1063 - [Producer clientId=producer-1] Bootstrap broker localhost:9093 (id: -1 rack: null) disconnected
16:17:42.434 [org.springframework.kafka.KafkaListenerEndpointContainer#0-0-C-1] WARN  o.apache.kafka.clients.NetworkClient:1063 - [Consumer clientId=consumer-sales.quotation.quotation.receiver-1, groupId=sales.quotation.quotation.receiver] Bootstrap broker localhost:9093 (id: -1 rack: null) disconnected
16:17:42.542 [kafka-producer-network-thread | producer-1] INFO  o.apache.kafka.clients.NetworkClient:935 - [Producer clientId=producer-1] Node -1 disconnected.

配置遗漏问题及修复方法

当前配置存在3个核心问题,按影响优先级排序:

  1. 客户端未指定SSL安全协议
    生产者、消费者、Spring自动创建的Admin Client配置中,都缺少security.protocol=SSL配置。当前客户端默认用PLAINTEXT协议连接9093的SSL端口,TCP握手后协议不匹配直接被Broker断开,这是Node -1 disconnected报错的核心原因。
    需要在生产者、消费者的配置Map中补充:

    props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL");
    

    注意Spring自带的Kafka Admin客户端不会读取自定义的Producer/Consumer Factory配置,需要在application.properties中补充全局SSL配置,避免Admin Client连接失败:

    spring.kafka.properties.security.protocol=SSL
    spring.kafka.properties.ssl.keystore.location=path/to/keystore.jks
    spring.kafka.properties.ssl.keystore.password=password
    spring.kafka.properties.ssl.key.password=password
    spring.kafka.properties.ssl.truststore.location=path/to/truststore.jks
    spring.kafka.properties.ssl.truststore.password=password
    spring.kafka.properties.ssl.endpoint.identification.algorithm=
    

    其中ssl.endpoint.identification.algorithm配置为空即可,不要传空格字符串,部分JDK版本会把空格识别为非法算法名抛出异常。

  2. 证书路径配置无效
    代码中写的path\\to\\xxx.jks是占位符路径,需要替换为本地证书文件的实际绝对路径:Windows环境保留反斜杠,Linux/macOS环境改用正斜杠。路径错误会导致客户端加载证书失败,SSL握手直接中断。

  3. Broker端PEM证书配置不全
    wurstmeister/kafka镜像使用PEM格式证书时,仅配置keystore和truststore路径无法正确加载证书,需要在docker-compose.yml的broker环境变量中新增:

    KAFKA_SSL_KEYSTORE_CERTIFICATE_CHAIN: /etc/kafka/secrets/broker.crt
    KAFKA_SSL_KEYSTORE_KEY: /etc/kafka/secrets/broker.key
    KAFKA_SSL_TRUSTSTORE_CERTIFICATES: /etc/kafka/secrets/root.crt
    

    否则Broker的SSL端口无法正常完成握手流程。

额外注意:你开启了KAFKA_SSL_CLIENT_AUTH=required双向认证,必须保证生产者、消费者的keystore证书是由Broker truststore中存储的根CA签发的,否则SSL握手会被Broker拒绝。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 23:18:48