SpringBoot SSL连接Docker部署Kafka报API_VERSIONS请求取消问题
我搭建了一个集成简易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个核心问题,按影响优先级排序:
客户端未指定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版本会把空格识别为非法算法名抛出异常。证书路径配置无效
代码中写的path\\to\\xxx.jks是占位符路径,需要替换为本地证书文件的实际绝对路径:Windows环境保留反斜杠,Linux/macOS环境改用正斜杠。路径错误会导致客户端加载证书失败,SSL握手直接中断。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

