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

Spring Boot中Azure Service Bus会话接收代码如何改用Mono.usingWhen

代码改写方案

改写后完整代码

@Autowired
ServiceBusSessionReceiverAsyncClient apiMessageQueueIntegrator;
// ... 其他业务代码
Mono<ServiceBusReceiverAsyncClient> receiverMono = apiMessageQueueIntegrator.acceptSession(sessionid);
Disposable subscription = Mono.usingWhen(
        receiverMono,
        // 直接调用官方单条消息接收方法,返回Mono类型的单发布源
        receiver -> receiver.receiveMessage(),
        // 资源清理逻辑保持不变,会话使用完成后自动关闭
        receiver -> Mono.fromRunnable(receiver::close)
    )
    .subscribe(
        message -> {
            // 原有消息处理逻辑无需修改
            logger.info(String.format("Message received from queue. Session id: %s. Contents: %s%n",
                message.getSessionId(), message.getBody()));
            receivedMessage.setReceivedMessage(message);
            // 若无其他异步等待逻辑可直接删除CountDownLatch相关代码,Mono单条消费完成后会自动终止流
            timeoutCheck.countDown();
        },
        error -> {
            logger.info("Queue error occurred: " + error);
        }
    );

核心改动说明

  • 将原有的Flux.usingWhen替换为Mono.usingWhen,适配单条消息的发布场景
  • 资源消费逻辑从返回多消息流的receiver.receiveMessages()改为官方提供的单条消息接收方法receiver.receiveMessage(),该方法直接返回Mono<ServiceBusReceivedMessage>,天然只发射一次消息完成信号
  • 无需额外手动销毁订阅,Mono在消息消费完成/出现异常后都会自动触发后续的资源关闭逻辑,代码健壮性更高

内容的提问来源于stack exchange,提问作者Ziggy000

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 19:45:08