使用ClientSecretCredential向Azure Service Bus发消息遇连接重置错误
问题描述
尝试通过ClientSecretCredential使用Azure SDK向Azure Service Bus队列发送消息时,出现连接重置错误。代码无编译问题,依赖包为最新版本,改用HTTPS方式发送则完全正常。
代码示例
import com.azure.identity.ClientSecretCredential; import com.azure.identity.ClientSecretCredentialBuilder; import com.azure.messaging.servicebus.ServiceBusClientBuilder; import com.azure.messaging.servicebus.ServiceBusMessage; import com.azure.messaging.servicebus.ServiceBusSenderClient; public class App { /* variables removed for confidentiality*/ public void sendMessage() throws InterruptedException,Exception{ ServiceBusMessage guestCheckInEvent = new ServiceBusMessage(data); guestCheckInEvent.setContentType(contentType); guestCheckInEvent.setMessageId(messageId); System.setProperty("AZURE_CLIENT_ID", clientId); System.setProperty("AZURE_CLIENT_SECRET", clientSecret); System.setProperty("AZURE_TENANT_ID", tenantId); ClientSecretCredential credential = new ClientSecretCredentialBuilder() .clientId(clientId) .tenantId(tenantId) .clientSecret(clientSecret) .build(); ServiceBusSenderClient sender = new ServiceBusClientBuilder() .fullyQualifiedNamespace(nameSpace) .credential(credential) .sender() .queueName(queueName) .buildClient(); System.out.println("Sending message"); sender.sendMessage(guestCheckInEvent); System.out.println("Sent message"); sender.close(); } public static void main(String[] args) { App a=new App(); try { a.sendMessage(); } catch (Exception e) { e.printStackTrace(); } } }
错误信息
com.azure.messaging.servicebus.ServiceBusException: Retries exhausted: 3/3 at com.azure.messaging.servicebus.ServiceBusSenderAsyncClient.mapError(ServiceBusSenderAsyncClient.java:823) at reactor.core.publisher.Mono.lambda$onErrorMap$28(Mono.java:3773) at reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onError(FluxOnErrorResume.java:94) at reactor.core.publisher.MonoFlatMap$FlatMapMain.onError(MonoFlatMap.java:180) at reactor.core.publisher.MonoPeekTerminal$MonoTerminalPeekSubscriber.onError(MonoPeekTerminal.java:258) at reactor.core.publisher.SerializedSubscriber.onError(SerializedSubscriber.java:124) at reactor.core.publisher.FluxRetryWhen$RetryWhenMainSubscriber.whenError(FluxRetryWhen.java:225) at reactor.core.publisher.FluxRetryWhen$RetryWhenOtherSubscriber.onError(FluxRetryWhen.java:274) at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onError(FluxContextWrite.java:121) at reactor.core.publisher.FluxConcatMapNoPrefetch$FluxConcatMapNoPrefetchSubscriber.maybeOnError(FluxConcatMapNoPrefetch.java:326) at reactor.core.publisher.FluxConcatMapNoPrefetch$FluxConcatMapNoPrefetchSubscriber.onNext(FluxConcatMapNoPrefetch.java:211) at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onNext(FluxContextWrite.java:107) at reactor.core.publisher.SinkManyEmitterProcessor.drain(SinkManyEmitterProcessor.java:471) at reactor.core.publisher.SinkManyEmitterProcessor.tryEmitNext(SinkManyEmitterProcessor.java:269) at reactor.core.publisher.SinkManySerialized.tryEmitNext(SinkManySerialized.java:100) at reactor.core.publisher.InternalManySink.emitNext(InternalManySink.java:27) at reactor.core.publisher.FluxRetryWhen$RetryWhenMainSubscriber.onError(FluxRetryWhen.java:190) at reactor.core.publisher.SerializedSubscriber.onError(SerializedSubscriber.java:124) at reactor.core.publisher.SerializedSubscriber.onError(SerializedSubscriber.java:124) at reactor.core.publisher.FluxTimeout$TimeoutMainSubscriber.onError(FluxTimeout.java:219) at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onError(FluxPeekFuseable.java:234) at reactor.core.publisher.MonoFlatMap$FlatMapMain.secondError(MonoFlatMap.java:241) at reactor.core.publisher.MonoFlatMap$FlatMapInner.onError(MonoFlatMap.java:315) at reactor.core.publisher.MonoFlatMap$FlatMapMain.onError(MonoFlatMap.java:180) at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onError(FluxMapFuseable.java:142) at reactor.core.publisher.MonoFlatMap$FlatMapMain.onError(MonoFlatMap.java:180) at reactor.core.publisher.FluxMapFuseable$MapFuseableSubscriber.onError(FluxMapFuseable.java:142) at reactor.core.publisher.MonoPeekTerminal$MonoTerminalPeekSubscriber.onError(MonoPeekTerminal.java:258) at reactor.core.publisher.FluxHide$SuppressFuseableSubscriber.onError(FluxHide.java:142) at reactor.core.publisher.MonoIgnoreThen$ThenIgnoreMain.onError(MonoIgnoreThen.java:278) at reactor.core.publisher.SerializedSubscriber.onError(SerializedSubscriber.java:124) at reactor.core.publisher.SerializedSubscriber.onError(SerializedSubscriber.java:124) at reactor.core.publisher.FluxTimeout$TimeoutMainSubscriber.onError(FluxTimeout.java:219) at reactor.core.publisher.MonoNext$NextSubscriber.onError(MonoNext.java:93) at reactor.core.publisher.FluxFilterFuseable$FilterFuseableSubscriber.onError(FluxFilterFuseable.java:162) at reactor.core.publisher.FluxReplay$SizeBoundReplayBuffer.replayNormal(FluxReplay.java:865) at reactor.core.publisher.FluxReplay$SizeBoundReplayBuffer.replay(FluxReplay.java:965) at reactor.core.publisher.FluxReplay$ReplaySubscriber.onError(FluxReplay.java:1360) at reactor.core.publisher.FluxPeek$PeekSubscriber.onError(FluxPeek.java:222) at reactor.core.publisher.FluxOnErrorResume$ResumeSubscriber.onError(FluxOnErrorResume.java:106) at reactor.core.publisher.MonoIgnoreThen$ThenIgnoreMain.onError(MonoIgnoreThen.java:278) at reactor.core.publisher.MonoIgnoreThen$ThenIgnoreMain.subscribeNext(MonoIgnoreThen.java:231) at reactor.core.publisher.MonoIgnoreThen$ThenIgnoreMain.onComplete(MonoIgnoreThen.java:203) at reactor.core.publisher.SinkEmptyMulticast$VoidInner.complete(SinkEmptyMulticast.java:272) at reactor.core.publisher.SinkEmptyMulticast.tryEmitEmpty(SinkEmptyMulticast.java:86) at reactor.core.publisher.SinkEmptySerialized.tryEmitEmpty(SinkEmptySerialized.java:46) at reactor.core.publisher.InternalEmptySink.emitEmpty(InternalEmptySink.java:26) at com.azure.core.amqp.implementation.ReactorConnection.lambda$closeConnectionWork$35(ReactorConnection.java:540) at reactor.core.publisher.MonoRunnable.call(MonoRunnable.java:73) at reactor.core.publisher.MonoRunnable.call(MonoRunnable.java:32) at reactor.core.publisher.MonoIgnoreThen$ThenIgnoreMain.subscribeNext(MonoIgnoreThen.java:228) at reactor.core.publisher.MonoIgnoreThen$ThenIgnoreMain.onComplete(MonoIgnoreThen.java:203) at reactor.core.publisher.SinkEmptyMulticast$VoidInner.complete(SinkEmptyMulticast.java:272) at reactor.core.publisher.SinkEmptyMulticast.tryEmitEmpty(SinkEmptyMulticast.java:86) at reactor.core.publisher.SinkEmptySerialized.tryEmitEmpty(SinkEmptySerialized.java:46) at reactor.core.publisher.InternalEmptySink.emitEmpty(InternalEmptySink.java:26) at com.azure.core.amqp.implementation.ReactorExecutor.close(ReactorExecutor.java:188) at com.azure.core.amqp.implementation.ReactorExecutor.lambda$scheduleCompletePendingTasks$1(ReactorExecutor.java:173) at reactor.core.scheduler.SchedulerTask.call(SchedulerTask.java:68) at reactor.core.scheduler.SchedulerTask.call(SchedulerTask.java:28) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) Suppressed: java.lang.Exception: #block terminated with an error at reactor.core.publisher.BlockingSingleSubscriber.blockingGet(BlockingSingleSubscriber.java:139) at reactor.core.publisher.Mono.block(Mono.java:1734) at com.azure.messaging.servicebus.ServiceBusSenderClient.sendMessage(ServiceBusSenderClient.java:192) at App.sendMessage(App.java:44) at App.main(App.java:52) Caused by: reactor.core.Exceptions$RetryExhaustedException: Retries exhausted: 3/3 at reactor.core.Exceptions.retryExhausted(Exceptions.java:306) at reactor.util.retry.RetryBackoffSpec.lambda$static$0(RetryBackoffSpec.java:68) at reactor.util.retry.RetryBackoffSpec.lambda$null$4(RetryBackoffSpec.java:560) at reactor.core.publisher.FluxConcatMapNoPrefetch$FluxConcatMapNoPrefetchSubscriber.onNext(FluxConcatMapNoPrefetch.java:183) ... 55 more Caused by: com.azure.core.amqp.exception.AmqpException: Connection reset by peer, errorContext[NAMESPACE: <Namespace>. ERROR CONTEXT: N/A] at com.azure.core.amqp.implementation.ExceptionUtil.toException(ExceptionUtil.java:85) at com.azure.core.amqp.implementation.handler.ConnectionHandler.notifyErrorContext(ConnectionHandler.java:351) at com.azure.core.amqp.implementation.handler.ConnectionHandler.onTransportError(ConnectionHandler.java:253) at org.apache.qpid.proton.engine.BaseHandler.handle(BaseHandler.java:191) at org.apache.qpid.proton.engine.impl.EventImpl.dispatch(EventImpl.java:108) at org.apache.qpid.proton.reactor.impl.ReactorImpl.dispatch(ReactorImpl.java:324) at org.apache.qpid.proton.reactor.impl.ReactorImpl.process(ReactorImpl.java:291) at com.azure.core.amqp.implementation.ReactorExecutor.run(ReactorExecutor.java:91) ... 8 more
解决方案
核心错误为Connection reset by peer,结合HTTPS方式正常的表现,大概率是AMQP协议端口被网络限制,或默认传输方式不兼容当前网络环境,可按以下方案排查解决:
1. 强制使用AMQP over WebSockets
Azure Service Bus支持通过443端口(与HTTPS一致)的WebSocket传输AMQP,能绕过多数端口拦截。修改客户端构建代码,添加传输类型配置:
ServiceBusSenderClient sender = new ServiceBusClientBuilder() .fullyQualifiedNamespace(nameSpace) .credential(credential) .sender() .queueName(queueName) .transportType(ServiceBusTransportType.AMQP_WEB_SOCKETS) // 启用WebSocket传输 .buildClient();
2. 检查网络与防火墙规则
- 确认客户端网络允许出站访问Azure Service Bus的AMQP端口(5671)或WebSocket端口(443)
- 检查Azure Service Bus命名空间的防火墙设置,是否将客户端IP加入允许列表
- 若使用企业代理,需确保代理放行对应端口和协议的流量
3. 优化重试策略
默认重试次数可能不足以应对临时网络波动,可自定义重试配置:
ServiceBusSenderClient sender = new ServiceBusClientBuilder() .fullyQualifiedNamespace(nameSpace) .credential(credential) .sender() .queueName(queueName) .retryOptions(new ServiceBusRetryOptions() .setMaxRetries(5) .setDelay(Duration.ofSeconds(2)) .setMaxDelay(Duration.ofSeconds(10))) .buildClient();
4. 验证服务主体权限
确保使用的ClientSecretCredential对应的服务主体已被分配Azure Service Bus Data Sender角色,权限不足也可能引发连接异常。
内容的提问来源于stack exchange,提问作者user3725074
相关产品推荐
相关产品推荐

