Java无法接收Azure Service Bus Queue消息,求排查解决建议
问题原因与解决方案
核心问题分析
你的代码问题出在异步客户端的使用逻辑错误,以及未正确处理消息生命周期和异步操作的等待机制:
第一个版本
receiveMessages():CountDownLatch从未被触发countDown(),导致程序一直等待60秒超时后直接退出;同时开启了disableAutoComplete()却未手动完成消息,若之前有未处理的消息处于锁定状态,新的接收操作无法获取到消息。receiveMessages()是持续接收消息的异步流,若队列消息处于锁定状态(比如之前测试未完成消息),即便peek能看到,也无法被正常接收。
第二个版本
receiveMessagesV2():- 循环中重复订阅
receiveMessages()是错误用法,异步订阅会立即返回,你没有等待回调执行就进入下一次循环,最后finally块直接关闭receiver,导致异步回调还没处理消息就被终止。 - 在异步回调中直接调用
complete却未处理其异步返回结果,容易出现操作未完成就关闭客户端的情况。
- 循环中重复订阅
修正方案
方案1:修复异步客户端代码
调整逻辑,在消息回调中完成消息并触发CountDownLatch,同时处理订阅的完成与错误信号:
public void receiveMessages() throws Exception { AtomicBoolean sampleSuccessful = new AtomicBoolean(true); CountDownLatch countdownLatch = new CountDownLatch(1); try (ServiceBusReceiverAsyncClient receiver = new ServiceBusClientBuilder() .connectionString(connectionString) .receiver() .queueName(queueName) .maxAutoLockRenewDuration(Duration.ofMinutes(1)) .disableAutoComplete() .buildAsyncClient()) { Disposable subscription = receiver.receiveMessages() .subscribe( message -> { System.out.printf("Sequence #: %s. Contents: %s%n", message.getSequenceNumber(), message.getBody()); // 手动完成消息,避免消息锁定到期后重新入队 receiver.complete(message) .doOnSuccess(v -> countdownLatch.countDown()) .doOnError(error -> { System.err.println("Failed to complete message: " + error); sampleSuccessful.set(false); countdownLatch.countDown(); }) .subscribe(); }, error -> { System.err.println("Error occurred while receiving message: " + error); sampleSuccessful.set(false); countdownLatch.countDown(); }, () -> countdownLatch.countDown()); if (!countdownLatch.await(60, TimeUnit.SECONDS)) { System.out.println("No message received within timeout period."); } subscription.dispose(); } catch (Exception ex) { ex.printStackTrace(); throw ex; } }
方案2:改用同步客户端(更适合新手)
异步客户端需处理复杂的响应式流逻辑,同步客户端ServiceBusReceiverClient逻辑更直观,易排查问题:
public void receiveMessagesSync() throws Exception { try (ServiceBusReceiverClient receiver = new ServiceBusClientBuilder() .connectionString(connectionString) .receiver() .queueName(queueName) .receiveMode(ServiceBusReceiveMode.PEEK_LOCK) .buildClient()) { // 批量获取消息,最多10条,等待5秒超时 Iterable<ServiceBusReceivedMessage> messages = receiver.receiveMessages(10, Duration.ofSeconds(5)); for (ServiceBusReceivedMessage message : messages) { System.out.printf("Sequence #: %s. Contents: %s%n", message.getSequenceNumber(), message.getBody()); // 手动完成消息 receiver.complete(message); } } catch (Exception ex) { ex.printStackTrace(); throw ex; } }
额外排查点
- 检查队列消息状态:在Azure门户的Service Bus队列页面查看“活动消息数”,若消息被锁定,等待锁到期(默认30秒)后再测试,或手动解锁消息。
- 确认权限有效性:再次检查连接字符串对应的共享访问策略,确保包含
Listen权限。
内容的提问来源于stack exchange,提问作者itsactuallysats
相关产品推荐
相关产品推荐

