Spring Boot 3.1.X及以上版本Kafka客户端连接问题求助
Spring Boot 3.1.X及以上版本升级后Kafka连接失败,回退至3.0.X恢复正常
近期将Spring Boot服务升级到3.1.X版本后,遇到Kafka连接问题,服务无法正常连接Confluent Cloud上的Kafka集群,持续输出如下日志:
2024-01-03T06:18:46.313456838Z 2024-01-03T06:18:46.313Z WARN 1 --- [servicerequestorms] [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-Dev_Phoenix_ServiceRequestor_CG-1, groupId=Dev_Phoenix_ServiceRequestor_CG] Bootstrap broker localhost:9092 (id: -1 rack: null) disconnected Wed, Jan 3 2024 11:48:47 am 2024-01-03T06:18:47.403Z INFO 1 --- [servicerequestorms] [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-Dev_Phoenix_ServiceRequestor_CG-1, groupId=Dev_Phoenix_ServiceRequestor_CG] Node -1 disconnected. 2024-01-03T06:18:47.403489607Z 2024-01-03T06:18:47.403Z WARN 1 --- [servicerequestorms] [ntainer#0-0-C-1] org.apache.kafka.clients.NetworkClient : [Consumer clientId=consumer-Dev_Phoenix_ServiceRequestor_CG-1, groupId=Dev_Phoenix_ServiceRequestor_CG] Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.
回退到Spring Boot 3.0.X版本时一切正常。项目使用spring-kafka 3.1.1,尝试升级到最新的Spring Boot 3.2.1后问题依旧,仅在3.1.X及以上版本出现该问题。
我的Spring Kafka配置(application.yml):
kafka: consumer: key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: io.confluent.kafka.serializers.KafkaJsonDeserializer group-id: Dev_Phoenix_ServiceRequestor_CG producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: io.confluent.kafka.serializers.KafkaAvroSerializer properties: ssl.endpoint.identification.algorithm: https sasl.mechanism: PLAIN bootstrap.servers: xyz.eastus.azure.confluent.cloud:9092 sasl.jaas.config: org.apache.kafka.common.security.plain.PlainLoginModule required username='username' password='password'; security.protocol: SASL_SSL schema.registry.url: https://xyz.westus2.azure.confluent.cloud basic.auth.credentials.source: USER_INFO schema.registry.basic.auth.user.info: M2T3WWJYIPFI5NBO:qSZikyV70b1h7tEHtmenmLTje7+aeUi6h+ZwRZ6scwLYUzzo9EhI70bk8OuWEH4Y #specific.avro.reader: true json.value.type: com.xxx.phoenix.events.DomainEvent
KafkaAdmin配置类:
@Configuration public class DomainEventsKafkaTopicConfig { @Value("${spring.kafka.properties.bootstrap.servers}") private String bootstrapAddress; @Value("${spring.kafka.properties.sasl.jaas.config}") private String phoenixClusterJAASConfig; @Value("${spring.kafka.properties.sasl.mechanism:PLAIN}") private String phoenixClusterSASLMechanism; @Value("${spring.kafka.properties.security.protocol:SASL_SSL}") private String phoenixClusterSecurityProtocol; @Value("${phoenix.domain.event.topic.name}") private String domainEventTopic; @Value("${phoenix.domain.event.topic.partitions:3}") private int domainEventTopicPartitions; @Value("${phoenix.domain.event.topic.replicationFactor:3}") private int domainEventTopicReplicationFactor; public DomainEventsKafkaTopicConfig() { } @Bean public KafkaAdmin kafkaAdmin() { Map<String, Object> configs = new HashMap(); configs.put("bootstrap.servers", this.bootstrapAddress); configs.put("security.protocol", this.phoenixClusterSecurityProtocol); configs.put("sasl.mechanism", this.phoenixClusterSASLMechanism); configs.put("sasl.jaas.config", this.phoenixClusterJAASConfig); return new KafkaAdmin(configs); } @Bean public NewTopic topic1() { return new NewTopic(this.domainEventTopic, this.domainEventTopicPartitions, (short)this.domainEventTopicReplicationFactor); } }
目前无法完成版本升级,恳请各位帮忙排查问题。
内容的提问来源于stack exchange,提问作者Arpan
相关产品推荐
相关产品推荐

