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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 23:32:07