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

Java无法接收Azure Service Bus Queue消息,求排查解决建议

问题原因与解决方案

核心问题分析

你的代码问题出在异步客户端的使用逻辑错误,以及未正确处理消息生命周期和异步操作的等待机制:

  1. 第一个版本receiveMessages():

    • CountDownLatch从未被触发countDown(),导致程序一直等待60秒超时后直接退出;同时开启了disableAutoComplete()却未手动完成消息,若之前有未处理的消息处于锁定状态,新的接收操作无法获取到消息。
    • receiveMessages()是持续接收消息的异步流,若队列消息处于锁定状态(比如之前测试未完成消息),即便peek能看到,也无法被正常接收。
  2. 第二个版本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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 17:41:03