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

使用AMQP 1.0协议(ProtonJ2)连接RabbitMQ持久化队列失败问题及解决

问题:Apache ProtonJ2 AMQP 1.0消费者连接RabbitMQ持久化队列失败

连接时抛出以下错误:

org.apache.qpid.protonj2.client.exceptions.ClientSessionRemotelyClosedException: PRECONDITION_FAILED - inequivalent arg 'durable' for queue 'test_queue_durable' in vhost '/': received 'false' but current is 'true' [condition = amqp:precondition-failed]

原示例代码

public void connectAmqp() throws Throwable {
    final String serverHost = "localhost";
    final int serverPort = 5672;
    final String address = "test_queue_durable";
    final Client client = Client.create();

    final ConnectionOptions options = new ConnectionOptions().user("admin").password("admin");

    try{
        Connection connection = client.connect(serverHost, serverPort, options);
        
        Receiver receiver = connection.openReceiver(address);
        
        for (int i = 0; i < 100; ++i) {
            Delivery delivery = receiver.receive();
            System.out.println(delivery.message().body().getClass());
            System.out.println("*-*-*-*          " + new String((byte[])delivery.message().body()));
        }
    }catch (Exception e) {
        e.printStackTrace();
    }
}

关键背景信息

  • 队列已在RabbitMQ中预先声明
  • 队列配置为持久化以避免消息丢失
  • 移除队列持久化属性后消费者可正常连接,但这不符合需求
  • 必须使用AMQP 1.0协议连接
  • 使用的客户端库为Apache Qpid Proton J2

解决方案

问题根源是消费者未匹配队列的持久化配置,需在ReceiverOptions的Source选项中设置对应持久化模式。修改后的代码如下:

public void connectAmqp() throws Throwable {
    final String serverHost = "localhost";
    final int serverPort = 5672;
    final String address = "test_queue_durable";
    final Client client = Client.create();

    try{
        ConnectionOptions options = new ConnectionOptions().user("user").password("pa$$w0rd");
        Connection connection = client.connect(serverHost, serverPort, options);
        /**
         * 消费者连接持久化队列的配置
         * 在ReceiverOptions的Source选项中设置持久化属性
         */
        ReceiverOptions ro = new ReceiverOptions();
        ro.sourceOptions().durabilityMode(DurabilityMode.CONFIGURATION);
        
        Receiver receiver = connection.openReceiver(address, ro);
        
        for (int i = 0; i < 100; ++i) {
            Delivery delivery = receiver.receive();
            System.out.println(delivery.message().body().getClass());
            System.out.println("*-*-*-*          " + new String((byte[])delivery.message().body()));
        }
    }catch (Exception e) {
        e.printStackTrace();
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 10:07:40