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

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("-----&gt; FAIL !");
    }

    @Override
    public void onSuccess(Confirm confirm)
    {
        System.err.println("-----&gt; SUCCESS !");
        System.err.println("-----&gt; " + confirm.isAck()//
            + "\n" + lo_data.getReturned() ///
            + "\n" + lo_data.getReturnedMessage() //
        );
    }
});

输出结果

发送部分输出

-----&gt; SUCCESS !

-----&gt; 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 20:23:16