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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 20:42:31