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

