使用 RocketMQ 实现消息的顺序消费
创始人
2025-01-08 00:32:55
0

在分布式系统中,保持消息的顺序性是一个常见且重要的问题。RocketMQ 提供了一种有效的方式来确保消息的顺序消费。本文将通过代码示例,介绍如何使用 RocketMQ 实现消息的顺序生产和消费。

环境准备

在开始之前,请确保您已经配置好 RocketMQ 环境,并且在 MqConstant 类中定义了 RocketMQ 的 NameServer 地址。

顺序消息的生产

首先,我们需要编写生产者代码来发送顺序消息。我们会创建两个示例,一个简单的顺序生产示例,另一个则是基于业务逻辑(如订单流程)的顺序生产示例。

简单的顺序生产者

package com.takumilove.demo;  import com.takumilove.constant.MqConstant; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.common.message.Message; import org.junit.Test;  public class FOrderlyTest {      @Test     public void orderlyProducer() throws Exception {         DefaultMQProducer producer = new DefaultMQProducer("orderly-producer-group");         producer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);         producer.start();         for (int i = 0; i < 10; i++) {             Message message = new Message("orderlyTopic", ("我是第" + i + "个消息").getBytes());             producer.send(message);         }         producer.shutdown();         System.out.println("发送完毕:");     } } 

基于业务逻辑的顺序生产者

在这个示例中,我们假设有一个 Order 类表示订单,订单包含了 idorderNumberpricedatestatus 等信息。

package com.takumilove.demo;  import com.takumilove.constant.MqConstant; import com.takumilove.domain.Order; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.client.producer.MessageQueueSelector; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.message.MessageQueue; import org.junit.Test;  import java.util.Arrays; import java.util.Date; import java.util.List;  public class FOrderlyTest {      @Test     public void orderlyProducer() throws Exception {         DefaultMQProducer producer = new DefaultMQProducer("orderly-producer-group");         producer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);         producer.start();         List orderList = Arrays.asList(                 new Order(1, 111, 59D, new Date(), "下订单"),                 new Order(2, 111, 59D, new Date(), "物流"),                 new Order(3, 111, 59D, new Date(), "签收"),                 new Order(4, 112, 89D, new Date(), "下订单"),                 new Order(5, 112, 89D, new Date(), "物流"),                 new Order(6, 112, 89D, new Date(), "拒收")         );         orderList.forEach(order -> {             Message message = new Message("orderlyTopic", order.toString().getBytes());             try {                 producer.send(message, new MessageQueueSelector() {                     @Override                     public MessageQueue select(List list, Message message, Object o) {                         int queueNumber = list.size();                         Integer i = (Integer) o;                         return list.get(i % queueNumber);                     }                 }, order.getOrderNumber());             } catch (Exception e) {                 System.out.println("发送失败" + e.getMessage());             }         });         producer.shutdown();         System.out.println("发送完毕:");     } } 

顺序消息的消费

接下来,我们编写消费者代码来消费这些顺序消息。我们将分别展示简单顺序消费者和基于业务逻辑的顺序消费者。

简单的顺序消费者

package com.takumilove.demo;  import com.takumilove.constant.MqConstant; import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.common.message.MessageExt; import org.junit.Test;  import java.util.List;  public class FOrderlyTest {      @Test     public void orderlyConsumer() throws Exception {         DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("orderly-consumer-group");         consumer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);         consumer.subscribe("orderlyTopic", "*");         consumer.registerMessageListener(new MessageListenerConcurrently() {             @Override             public ConsumeConcurrentlyStatus consumeMessage(List list,                                                             ConsumeConcurrentlyContext consumeConcurrentlyContext) {                 System.out.println(new String(list.get(0).getBody()));                 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;             }         });         consumer.start();         System.in.read();     } } 

基于业务逻辑的顺序消费者

package com.takumilove.demo;  import com.takumilove.constant.MqConstant; import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly; import org.apache.rocketmq.common.message.MessageExt; import org.junit.Test;  import java.util.List;  public class FOrderlyTest {      @Test     public void orderlyConsumer() throws Exception {         DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("orderly-consumer-group");         consumer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);         consumer.subscribe("orderlyTopic", "*");         consumer.registerMessageListener(new MessageListenerOrderly() {             @Override             public ConsumeOrderlyStatus consumeMessage(List list,                                                        ConsumeOrderlyContext consumeOrderlyContext) {                 MessageExt messageExt = list.get(0);                 System.out.println(new String(messageExt.getBody()));                 return ConsumeOrderlyStatus.SUCCESS;             }         });         consumer.start();         System.in.read();     } } 

总结

通过以上示例,我们展示了如何使用 RocketMQ 实现消息的顺序生产和消费。无论是简单的消息还是基于业务逻辑的消息,都可以通过 RocketMQ 提供的顺序消费机制来保证消息的有序性。这对于订单系统等需要严格顺序的场景尤为重要。

相关内容

热门资讯

最终!德普之星透视免费,微信小... 最终!德普之星透视免费,微信小程序雀神,好像真的是有挂(有挂攻略)1、下载好微信小程序雀神脚本下载之...
黑科技技巧!wepoker免费... 黑科技技巧!wepoker免费脚本,jj斗地主外开挂,都是是真的有挂(讲解有挂)1、每一步都需要思考...
突发!pokemomo辅助软件... 突发!pokemomo辅助软件,winner辅助软件,切实是真的有挂(有挂教程)1、实时winner...
更值得关注的是!wpk真吗,潮... 更值得关注的是!wpk真吗,潮友会鱼虾蟹塞子概率计算方式,原来真的有挂(有挂细节)一、潮友会鱼虾蟹塞...
突发!aapoker怎么选牌,... 突发!aapoker怎么选牌,雀友会广东潮汕辅助脚本,总是真的是有挂(真的有挂)该软件可以轻松地帮助...
截至目前!wpk软件是正规的吗... 截至目前!wpk软件是正规的吗,宝宝浙江游戏免费开挂,一直是真的有挂(有挂解密)宝宝浙江游戏免费开挂...
现有关情况通报如下!哈糖大菠萝... 现有关情况通报如下!哈糖大菠萝万能挂,蜀山四川麻亲友房祈福,果然存在有挂(发现有挂)1、游戏颠覆性的...
黑科技辅助挂!wepoker开... 黑科技辅助挂!wepoker开脚本视频,多乐小程序游戏破解器,确实确实有挂(证实有挂)黑科技辅助挂!...
2026版复盘!wepoker... 2026版复盘!wepoker俱乐部辅助,财神破解版全自动脚本,切实真的有挂(存在有挂)1、下载好财...
突发!we-poker辅助器,... 突发!we-poker辅助器,枫叶辅助器,一贯真的是有挂(揭秘有挂)枫叶辅助器透视方法中分为三种模型...