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

Spring Boot中Reactor Kafka连接Azure Event Hubs出现Broker断开问题

问题根因

你当前的Reactor Kafka消费者缺少Azure Event Hubs Kafka兼容端点强制要求的TLS加密和SASL身份认证配置,从日志输出的ConsumerConfig可以看到,当前使用的是PLAINTEXT明文协议,SASL相关配置为空,因此会触发连接反复断开的问题。之前使用Spring Cloud Stream绑定时框架会自动注入相关认证配置,手动构造Reactor Kafka Receiver时需要自行补充。

修复方案

在消费者配置Map中补充以下4项核心配置即可:

  • 安全协议指定为SASL_SSL,适配Event Hubs的加密要求
  • SASL认证机制指定为PLAIN
  • 配置JAAS认证信息,填入Event Hubs命名空间的连接字符串
  • 开启SSL端点校验,避免证书认证失败

修改后的完整代码如下:

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.kafka.receiver.KafkaReceiver;
import reactor.kafka.receiver.ReceiverOptions;
import reactor.kafka.receiver.ReceiverRecord;

import java.util.Collections;
import java.util.HashMap;
import java.util.Map;

@Slf4j
@Service
public class StreamConsumer {

    public Flux<Object> consumeMessages() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "mynamespacename.servicebus.windows.net:9093");
        props.put(ConsumerConfig.CLIENT_ID_CONFIG, "test-client");
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "$Default");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, org.springframework.kafka.support.serializer.JsonDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, org.apache.kafka.common.serialization.StringDeserializer.class);

        // 新增Azure Event Hubs必要认证配置
        props.put("security.protocol", "SASL_SSL");
        props.put(SaslConfigs.SASL_MECHANISM, "PLAIN");
        // 注意password替换为你自己的Event Hubs命名空间连接字符串,格式为Endpoint=sb://xxx.servicebus.windows.net/;SharedAccessKeyName=xxx;SharedAccessKey=xxx
        props.put(SaslConfigs.SASL_JAAS_CONFIG, "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"你的Event Hubs连接字符串\";");
        props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, "https");

        ReceiverOptions<String, Object> receiverOptions = ReceiverOptions.create(props);

        ReceiverOptions<String, Object> options = receiverOptions.subscription(Collections.singleton(KafkaConstants.KAFKA_TOPIC))
                .addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions))
                .addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions));
        Flux<ReceiverRecord<String, Object>> kafkaFlux = KafkaReceiver.create(options).receive();
        return kafkaFlux.map(x -> "Test");
    }

}
额外检查项
  • 确认使用的Event Hubs连接字符串来自命名空间级别共享访问策略,且具备监听(Listen)权限
  • 若Event Hubs开启了IP白限制,确认当前服务的出口IP已加入允许列表
  • 若使用专用集群,确认9093端口的网络连通性正常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 06:18:03