Spring Boot连接Azure Event Hubs超时,无法接收消息求助
问题:Spring Boot连接Azure Event Hubs超时,无法接收消息
应用代码
@SpringBootApplication public class VALogAnalyticsApplication implements CommandLineRunner { private static final Sinks.Many<Message<String>> many = Sinks.many().unicast().onBackpressureBuffer(); public static void main(String[] args) { SpringApplication.run(VALogAnalyticsApplication.class, args); } @Override public void run(String... args) { log.info("Going to add message {} to sendMessage." + "Hello World"); many.emitNext(MessageBuilder.withPayload("Hello World").build(), Sinks.EmitFailureHandler.FAIL_FAST); } @Bean public Consumer<Message<String>> consume() { return message->{ Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER); log.info("New message received: '{}', partition key: {}, sequence number: {}, offset: {}, enqueued " +"time: {}", message.getPayload(), message.getHeaders().get(EventHubsHeaders.PARTITION_KEY), message.getHeaders().get(EventHubsHeaders.SEQUENCE_NUMBER), message.getHeaders().get(EventHubsHeaders.OFFSET), message.getHeaders().get(EventHubsHeaders.ENQUEUED_TIME) ); checkpointer.success() .doOnSuccess(success->log.info("Message '{}' successfully checkpointed", message.getPayload())) .doOnError(error->log.error("Exception found", error)) .block(); }; } }
application.properties配置
spring.cloud.azure.eventhubs.namespace=<eventHubNameSpace> spring.cloud.azure.eventhubs.connection-string=Endpoint=sb://test.servicebus.windows.net /;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=adfsad spring.cloud.azure.eventhubs.processor.checkpoint-store.account-name=<storagAaccountName> spring.cloud.azure.eventhubs.processor.checkpoint-store.account-key=<storageAccessKey> spring.cloud.azure.eventhubs.processor.checkpoint-store.container-name=<containerName> spring.cloud.stream.bindings.consume-in-0.destination=<eventHubName> spring.cloud.stream.bindings.consume-in-0.group=<consumerGroup> spring.cloud.stream.bindings.supply-out-0.destination=<eventHubName> spring.cloud.stream.eventhubs.bindings.consume-in-0.consumer.checkpoint.mode=MANUAL spring.cloud.function.definition=consume;supply; spring.cloud.stream.poller.initial-delay=0 spring.cloud.stream.poller.fixed-delay=1000
错误日志
2023-08-22 16:51:20.586 INFO 26256 --- [ctor-executor-1] c.a.c.a.i.ReactorDispatcher : {"az.sdk.message":"Reactor selectable is being disposed.","connectionId":"MF_89ee14_1692703254531"} 2023-08-22 16:51:20.587 INFO 26256 --- [ctor-executor-1] c.a.c.a.i.ReactorConnection : {"az.sdk.message":"onConnectionShutdown. Shutting down.","connectionId":"MF_89ee14_1692703254531","isTransient":false,"isInitiatedByClient":false,"shutdownMessage":"connectionId[MF_89ee14_1692703254531] Reactor selectable is disposed.","namespace":"test.servicebus.windows.net"} 2023-08-22 16:51:20.639 ERROR 26256 --- [ctor-executor-1] reactor.core.publisher.Operators : Operator called default onErrorDropped reactor.core.Exceptions$ErrorCallbackNotImplemented: com.azure.core.amqp.exception.AmqpException: Connection timed out: no further information, errorContext[NAMESPACE: test.servicebus.windows.net. ERROR CONTEXT: N/A] Caused by: com.azure.core.amqp.exception.AmqpException: Connection timed out: no further information, errorContext[NAMESPACE: test.servicebus.windows.net. ERROR CONTEXT: N/A] at com.azure.core.amqp.implementation.ExceptionUtil.toException(ExceptionUtil.java:85) ~[azure-core-amqp-2.8.7.jar:2.8.7] at com.azure.core.amqp.implementation.handler.ConnectionHandler.notifyErrorContext(ConnectionHandler.java:351) ~[azure-core-amqp-2.8.7.jar:2.8.7] at com.azure.core.amqp.implementation.handler.ConnectionHandler.onTransportError(ConnectionHandler.java:253) ~[azure-core-amqp-2.8.7.jar:2.8.7] at org.apache.qpid.proton.engine.BaseHandler.handle(BaseHandler.java:191) ~[proton-j-0.33.8.jar:na] at org.apache.qpid.proton.engine.impl.EventImpl.dispatch(EventImpl.java:108) ~[proton-j-0.33.8.jar:na] at org.apache.qpid.proton.reactor.impl.ReactorImpl.dispatch(ReactorImpl.java:324) ~[proton-j-0.33.8.jar:na] at org.apache.qpid.proton.reactor.impl.ReactorImpl.process(ReactorImpl.java:291) ~[proton-j-0.33.8.jar:na] at com.azure.core.amqp.implementation.ReactorExecutor.run(ReactorExecutor.java:91) ~[azure-core-amqp-2.8.7.jar:2.8.7] at reactor.core.scheduler.SchedulerTask.call(SchedulerTask.java:68) ~[reactor-core-3.4.6.jar:3.4.6] at reactor.core.scheduler.SchedulerTask.call(SchedulerTask.java:28) ~[reactor-core-3.4.6.jar:3.4.6] at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[na:na] at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304) ~[na:na] at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) ~[na:na] Caused by: com.azure.core.amqp.exception.AmqpException: Connection timed out: no further information, errorContext[NAMESPACE: test.servicebus.windows.net. ERROR CONTEXT: N/A]
排查与修正步骤
1. 修复连接字符串格式错误
配置中的spring.cloud.azure.eventhubs.connection-string存在格式问题:Endpoint的URL里有多余空格sb://test.servicebus.windows.net /,需改为sb://test.servicebus.windows.net/,空格会导致寻址失败直接引发超时。
2. 检查网络与防火墙规则
- 确认应用所在环境可访问Azure Event Hubs的端口:AMQP用5671端口,HTTPS用443端口,企业内网需确保防火墙未拦截这些端口。
- 登录Azure门户检查Event Hubs命名空间的防火墙设置,确认允许应用的公网IP访问。
3. 验证权限与资源配置
- 确认连接字符串使用的SharedAccessKey具备Listen权限(消费消息必需),生产环境建议使用最小权限的专用密钥,而非RootManageSharedAccessKey。
- 检查检查点存储的存储账户、容器名称是否正确,存储密钥是否有效;存储账户需与Event Hubs命名空间网络可达,跨区域部署可能引发超时。
4. 代码优化
- 日志拼接错误:
log.info("Going to add message {} to sendMessage." + "Hello World");占位符未对应参数,改为log.info("Going to add message {} to sendMessage.", "Hello World");。 - 消费方法中
checkpointer.success().block()会阻塞线程,Reactor环境建议改为非阻塞处理:
checkpointer.success() .doOnSuccess(success->log.info("Message '{}' successfully checkpointed", message.getPayload())) .doOnError(error->log.error("Exception found", error)) .subscribe();
5. 依赖版本检查
确保spring-cloud-azure-eventhubs与Azure SDK版本匹配,避免版本冲突导致连接问题,建议使用官方推荐的版本组合(如spring-cloud-azure 4.x搭配azure-core-amqp 2.10+)。
内容的提问来源于stack exchange,提问作者Prem Kumar
相关产品推荐
相关产品推荐

