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

Java中IBM MQ同步收发消息:实现发订单后等发货回复

实现JMS同步请求-回复流程的两种方案

方案一:使用Spring JMS内置的sendAndReceive方法

Spring JMS的JmsTemplate已经封装好了同步请求-回复逻辑,直接调用sendAndReceive就能完成「发送消息→等待回复→返回结果」的流程,无需手动处理监听和等待逻辑,代码最简洁。

修改后的placeOrder方法代码

// 将队列名称和依赖定义为类成员,不要放在方法内部
private static final String SEND_QUEUE = "order_queue";
private static final String REPLY_QUEUE = "shipment_queue";
private final JmsTemplate jmsTemplate;
private final OurConverter ourConverter; // 你的消息转换器

public Shipment placeOrder(Order order) throws JMSException {
    // 发送消息并同步阻塞等待回复
    Message replyMessage = jmsTemplate.sendAndReceive(SEND_QUEUE, session -> {
        Message jmsMsg = ourConverter.toMessage(order, session);
        jmsMsg.setJMSExpiration(5 * Constants.MINUTE);
        jmsMsg.setJMSDeliveryMode(DeliveryMode.NON_PERSISTENT);
        jmsMsg.setJMSReplyTo(session.createQueue(REPLY_QUEUE));
        return jmsMsg;
    });

    // 转换回复消息为Shipment返回
    if (replyMessage != null) {
        ReplyData replyData = ourConverter.replyFromMessage(replyMessage);
        return convertToShipment(replyData);
    } else {
        throw new RuntimeException("未收到订单回复,请求超时");
    }
}

// 辅助方法:根据业务逻辑将ReplyData转换为Shipment
private Shipment convertToShipment(ReplyData replyData) {
    Shipment shipment = new Shipment();
    // 填充Shipment字段,比如replyData.getShipmentId()等
    return shipment;
}

关键点说明

  • sendAndReceive会自动生成唯一JMSCorrelationID,并阻塞当前线程直到收到对应回复或超时
  • 无需保留原异步@JmsListener方法,Spring内部会处理监听逻辑
  • 可通过jmsTemplate.setReceiveTimeout()全局设置超时时间,或在方法中捕获超时异常

方案二:手动基于CorrelationID+Future实现(自定义场景)

如果需要对回复逻辑做更灵活控制(比如自定义缓存、多线程隔离),可以手动通过CorrelationID关联请求和回复,结合CompletableFuture实现同步等待。

步骤1:添加全局缓存和依赖

在JMS处理类中添加线程安全的缓存,存储CorrelationID对应的Future:

private static final String SEND_QUEUE = "order_queue";
private static final String REPLY_QUEUE = "shipment_queue";
private final JmsTemplate jmsTemplate;
private final OurConverter ourConverter;
// 线程安全缓存,存CorrelationID -> 对应的CompletableFuture
private final ConcurrentHashMap<String, CompletableFuture<ReplyData>> replyFutureCache = new ConcurrentHashMap<>();

步骤2:修改placeOrder方法发送消息并等待回复

public Shipment placeOrder(Order order) throws InterruptedException, ExecutionException, TimeoutException {
    String correlationId = UUID.randomUUID().toString();
    CompletableFuture<ReplyData> replyFuture = new CompletableFuture<>();
    replyFutureCache.put(correlationId, replyFuture);

    try {
        jmsTemplate.send(SEND_QUEUE, session -> {
            Message jmsMsg = ourConverter.toMessage(order, session);
            jmsMsg.setJMSCorrelationID(correlationId);
            jmsMsg.setJMSExpiration(5 * Constants.MINUTE);
            jmsMsg.setJMSDeliveryMode(DeliveryMode.NON_PERSISTENT);
            jmsMsg.setJMSReplyTo(session.createQueue(REPLY_QUEUE));
            return jmsMsg;
        });

        // 阻塞等待回复,设置和消息过期一致的超时时间
        ReplyData replyData = replyFuture.get(5, TimeUnit.MINUTES);
        return convertToShipment(replyData);
    } catch (JMSException e) {
        throw new RuntimeException("发送订单消息失败", e);
    } finally {
        // 清理缓存,避免内存泄漏
        replyFutureCache.remove(correlationId);
    }
}

步骤3:修改@JmsListener方法处理回复并唤醒线程

@JmsListener(destination = "shipment_queue")
public void receiveReply(Message message) {
    try {
        String correlationId = message.getJMSCorrelationID();
        if (correlationId == null) {
            logger.warn("收到无CorrelationID的回复消息,忽略");
            return;
        }

        CompletableFuture<ReplyData> replyFuture = replyFutureCache.get(correlationId);
        if (replyFuture != null) {
            ReplyData replyData = ourConverter.replyFromMessage(message);
            replyFuture.complete(replyData);
        } else {
            logger.warn("收到未匹配请求的回复消息,CorrelationID: {}", correlationId);
        }
    } catch (JMSException e) {
        logger.error("处理回复消息失败", e);
    }
}

关键点说明

  • ConcurrentHashMap保证多线程环境下的缓存操作安全
  • finally块中移除缓存的Future,防止内存泄漏
  • 通过get(timeout, unit)设置超时,避免线程无限等待
  • 异常场景可调用replyFuture.completeExceptionally(e)将异常传递给等待线程

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 06:02:08