Spring AMQP Java中发送到不存在交换机失败时如何获取消息体
问题
在Java的Spring AMQP环境下,当消息发送到不存在的交换机时,希望获取消息体。目前遇到以下问题:
- 使用
confirmCallback时,能拿到correlation data ID、错误原因,但返回的消息为空 - 使用
CorrelationData的Future回调时,回调被标记为成功调用,无法获取有效错误信息 - 不想用数据库存储消息体和correlation data的UUID,希望在confirm回调中直接获取消息体
相关代码
连接工厂配置
CachingConnectionFactory connectionFactory = new CachingConnectionFactory(); connectionFactory.setUsername(as_userName); connectionFactory.setPassword(as_pwd); connectionFactory.setVirtualHost(PONOSEnvConst.QUEUE_VHOST); connectionFactory.setHost(PONOSEnvConst.URL_SERV_MQ); connectionFactory.setPort(PONOSEnvConst.PORT_SERV_MQ); connectionFactory.setPublisherReturns(true); connectionFactory.setPublisherConfirms(true);
RabbitTemplate配置
RabbitTemplate rabbitTemplate = new RabbitTemplate(ao_factory); rabbitTemplate.setMessageConverter(jsonMessageConverter()); rabbitTemplate.setReturnsCallback(ao_returned -> { System.err.println("Returned: " + ao_returned); }); rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> { System.err.println("##### confirm : ##### ack = " + ack // + "\n cause = " + cause // + "\n data = " + correlationData // + "\n returned = " + correlationData.getReturned() // + "\n######"); });
发送消息代码
CorrelationData lo_data = new CorrelationData(UUID.randomUUID().toString()); amqpTemplate.convertAndSend("non-existing-exchange", "", "Message body to an non-existing exchange", lo_data); lo_data.getFuture().addCallback(new ListenableFutureCallback<Confirm>() { @Override public void onFailure(Throwable throwable) { System.err.println("-----> FAIL !"); } @Override public void onSuccess(Confirm confirm) { System.err.println("-----> SUCCESS !"); System.err.println("-----> " + confirm.isAck()// + "\n" + lo_data.getReturned() /// + "\n" + lo_data.getReturnedMessage() // ); } });
输出结果
发送部分输出
-----> SUCCESS ! -----> false null null
Confirm回调部分输出
##### confirm : ##### ack = false cause = channel error; protocol method: #method<channel.close>(reply-code=404, reply-text=NOT_FOUND - no exchange 'non-existing-exchange' in vhost 'XXX', class-id=60, method-id=40) data = CorrelationData [id=4a9dd7a0-27d5-456d-9018-4e7384038730] returned = null ######
解决方案
可以通过自定义CorrelationData子类的方式,在发送消息时将消息体或转换后的Message对象存入其中,这样在confirm回调中就能直接获取到消息内容,无需数据库存储。
步骤1:自定义CorrelationData子类
public class MessageCorrelationData extends CorrelationData { private final Object messageBody; private final Message rawMessage; public MessageCorrelationData(String id, Object messageBody, Message rawMessage) { super(id); this.messageBody = messageBody; this.rawMessage = rawMessage; } // Getter方法 public Object getMessageBody() { return messageBody; } public Message getRawMessage() { return rawMessage; } }
步骤2:发送消息时使用自定义CorrelationData
发送前先将消息转换为Message对象(或直接传入原消息体),创建自定义的CorrelationData实例:
// 原消息体 String messageBody = "Message body to an non-existing exchange"; // 转换为Message对象 Message message = rabbitTemplate.getMessageConverter().toMessage(messageBody, new MessageProperties()); // 创建自定义CorrelationData,传入ID、消息体、Message对象 MessageCorrelationData correlationData = new MessageCorrelationData(UUID.randomUUID().toString(), messageBody, message); // 发送消息 rabbitTemplate.send("non-existing-exchange", "", message, correlationData);
步骤3:在Confirm回调中获取消息体
修改confirm回调,将CorrelationData强转为自定义子类,直接获取消息内容:
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> { if (!ack) { if (correlationData instanceof MessageCorrelationData) { MessageCorrelationData msgCorrData = (MessageCorrelationData) correlationData; System.err.println("发送失败,消息体:" + msgCorrData.getMessageBody()); System.err.println("原始Message对象:" + msgCorrData.getRawMessage()); System.err.println("错误原因:" + cause); } } });
补充说明
- 发送到不存在的交换机属于特殊的“不可路由”场景,RabbitMQ会直接关闭通道,不会触发
returnsCallback,仅会触发confirm回调返回ack=false。 - 自定义CorrelationData的方式仅在内存中暂存消息体,直到confirm回调执行完成,无需持久化到数据库,符合需求。
内容的提问来源于stack exchange,提问作者Damien Schadeck
相关产品推荐
相关产品推荐

