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
相关产品推荐
相关产品推荐

