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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 17:19:52