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
相关产品推荐
相关产品推荐

