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

