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

Confluent Cloud Kafka集群连接失败求助(Spring Boot环境)

问题描述

在Confluent Cloud上创建Kafka集群后无法连接,运行生产者时出现以下错误:

[Producer clientId=producer-1] Node -1 disconnected.2023-06-05T06:09:20.826+05:30 INFO 25324 --- [ad | producer-1]org.apache.kafka.clients.NetworkClient : [ProducerclientId=producer-1] Cancelled in-flight API_VERSIONS request withcorrelation id 189 due to node -1 being disconnected (elapsed timesince creation: 253ms, elapsed time since send: 253ms, requesttimeout: 30000ms) 2023-06-05T06:09:20.827+05:30 WARN 25324 --- [ad |producer-1] org.apache.kafka.clients.NetworkClient : [ProducerclientId=producer-1] Bootstrap broker (id: -1 rack:null) disconnected

尝试创建新集群后问题依旧。使用Spring Boot连接,相关配置如下:

spring.kafka.properties.sasl.mechanism=PLAIN
spring.kafka.properties.bootstrap.servers=broker-address-here:9092
spring.kafka.properties.sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username='api-key-here' password='api-secret-here';
spring.kafka.properties.security.protocol=SASL_SSL
spring.kafka.properties.session.timeout.ms=45000

Spring Boot的Kafka配置Bean代码:

@Configuration
public class KafkaConfiguration {
  @Value("${spring.kafka.properties.bootstrap.servers}")
  private String bootStrapServer;
  @Bean
  public ProducerFactory<String, String> producerFactory() {
    return new DefaultKafkaProducerFactory<String, String>(Map.of(
            AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootStrapServer,
            AdminClientConfig.RETRIES_CONFIG, 0,
            ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432,
            ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class,
            ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class
    ));
  }
  @Bean
  public KafkaTemplate<String, String> kafkaTemplate() {
    return new KafkaTemplate<>(producerFactory());
  }
}

发送消息的控制器代码:

@RestController
@RequestMapping("/produce")
public final class LogProducer {
    @Autowired 
    private KafkaTemplate<String, String> kafkaTemplate;

    private static final String TOPIC = "logs";

    // Publish messages using the GetMapping
    @PostMapping("/logs/v1")
    public String publishMessage()
    {
        // Sending the message
        kafkaTemplate.send(TOPIC, "sample log message");
        return "Published Successfully";
    }
}

请问哪里操作出错了?


问题原因及解决方案

核心问题

你手动创建ProducerFactory时,只传入了基础配置参数,但没有包含Confluent Cloud要求的SASL/SSL安全认证参数。Confluent Cloud的Kafka集群必须通过SASL_SSL认证才能连接,而你的Map.of()里仅配置了bootstrap地址、重试、缓冲区和序列化器,缺失了sasl.mechanism、security.protocol、sasl.jaas.config这些关键安全配置,导致生产者无法完成认证,最终出现节点断开的错误。

解决方案

有两种修复方式:

方式一:复用Spring Boot自动加载的配置属性

直接使用Spring Boot自动配置的KafkaProperties构建ProducerFactory,自动读取application.properties里的所有Kafka配置(包括安全认证参数),无需手动拼接:

@Configuration
public class KafkaConfiguration {
  @Autowired
  private KafkaProperties kafkaProperties;

  @Bean
  public ProducerFactory<String, String> producerFactory() {
    return new DefaultKafkaProducerFactory<>(kafkaProperties.buildProducerProperties());
  }

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

方式二:手动添加所有必要的安全配置参数

如果坚持手动构建配置Map,需把application.properties里的安全认证参数全部补充进去:

@Configuration
public class KafkaConfiguration {
  @Value("${spring.kafka.properties.bootstrap.servers}")
  private String bootStrapServer;
  @Value("${spring.kafka.properties.sasl.mechanism}")
  private String saslMechanism;
  @Value("${spring.kafka.properties.sasl.jaas.config}")
  private String saslJaasConfig;
  @Value("${spring.kafka.properties.security.protocol}")
  private String securityProtocol;

  @Bean
  public ProducerFactory<String, String> producerFactory() {
    Map<String, Object> configs = new HashMap<>();
    configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootStrapServer);
    configs.put(AdminClientConfig.RETRIES_CONFIG, 0);
    configs.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
    configs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    // 添加安全认证参数
    configs.put("sasl.mechanism", saslMechanism);
    configs.put("sasl.jaas.config", saslJaasConfig);
    configs.put("security.protocol", securityProtocol);
    return new DefaultKafkaProducerFactory<>(configs);
  }

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

额外检查项

  1. 确认bootstrap.servers是Confluent Cloud控制台提供的正确地址,格式通常为pkc-xxxx.us-west-2.aws.confluent.cloud:9092
  2. 确认API Key和Secret是该集群对应的有效凭证,无拼写错误
  3. 检查网络是否能访问Confluent Cloud的9092端口,避免防火墙或代理拦截

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 13:30:34