RocketMq快速入门
2026/9/1 8:12:09 网站建设 项目流程
mqbroker.cmd -n 127.0.0.1:9876 autoCreateTopicEnable=true

开启确保 Broker 启动时就允许自动创建主题

对比维度Message(普通消息)MessageExt(扩展消息)
谁在用生产者(Producer)发送消息时使用消费者(Consumer)接收消息时使用
包含什么只有Topic、Tag、业务 Body(消息体)包含全部 Message 内容 + 系统元数据
用途只负责告诉 Broker “我要发什么内容”告诉消费者“你收到的这条消息,在服务器里的状态是什么”

正确的消费者管理方式

  1. 长期运行,保持单例:正确的模式是将消费者设计为长期运行的单例。在应用启动时创建并启动它,然后让它持续监听,应用关闭时才统一关闭。

  2. 优雅关闭,只做一次shutdown()应该只在你确定要彻底停止消息消费并释放资源时(例如应用正常下线),作为最后的清理步骤被调用一次。

  3. 避免在监听器中操作绝对不要MessageListenerconsumeMessage()方法内部去调用consumer.shutdown()。这会导致消费者在消息处理过程中被意外销毁。

RocketMQ的消息分发遵循“组内竞争,组间复制”原则:

  • 同组(Group相同)=竞争关系,一条消息只投递给组内一个消费者(负载均衡)。

  • 异组(Group不同)=复制关系,每条消息会被每个组独立消费一次(业务解耦)。

  1. 同组内的“争抢”不是随机抢单条,而是“抢队列”
    RocketMQ的实际分配单位是队列(MessageQueue),而不是单条消息。
    比如一个Topic有4个队列,组里有2个消费者,那么就是消费者A拿队列1和2,消费者B拿队列3和4(各自抢占一部分队列)。
    所以,同组内并不是每条消息都随机乱抢,而是先分好“责任田”(队列),再各自消费田里的消息

但是也可以在RocketMQ中同组使用广播模式即发送n条消息同一组下每个消费者都消费n条。

consumer.setMessageModel(MessageModel.BROADCASTING);

核心运行机制对比(底层原理)

对比维度同步消息 (Sync)异步消息 (Async)
线程状态发送线程阻塞(Blocking),等待网络IO和Broker响应。发送线程非阻塞(Non-blocking),发完即返回。
结果获取直接在send()方法的返回值里拿到SendResult通过实现SendCallback接口,在onSuccess()onException()里拿结果。
可靠性。能立即感知成功/失败,便于重试或回滚。较高。有回调机制,但需要处理好回调线程和主线程的上下文传递。

异步发送(send+Callback)后立即执行shutdown(),会因强制销毁连接和线程池,导致网络请求中断而报错。

单向消息:

// 单向发送:没有返回值,没有回调 producer.sendOneway(msg);

WaitStoreMsgOK是 RocketMQ 消息的一个属性,用于控制消息发送的可靠性

它的核心作用是:决定生产者发送消息时,是否需要等待 Broker(服务端)将消息成功写入磁盘(落盘)后,才返回成功响应。

  • WaitStoreMsgOK = true(默认):生产者发送消息后,会一直等待 Broker 的响应。Broker 收到消息后,会先将其写入磁盘(根据配置可能是同步或异步刷盘),只有在写入成功(或在一定超时时间内)后,才会向生产者返回SEND_OK状态。这保证了消息的高可靠性,但会带来额外的网络延迟,略微降低发送性能。

  • WaitStoreMsgOK = false:生产者发送消息后,Broker只要将消息写入内存(PageCache)即返回成功,无需等待落盘完成。这种方式发送性能更高、延迟更低,但可靠性降低,因为在 Broker 发生断电等故障时,未落盘的消息可能会丢失。

consumer.subscribe("topic6", MessageSelector.bySql("age>=18"));

  • 消费者模式:只有PUSH模式的消费者支持SQL过滤。

  • 订阅一致性:同一消费者组(ConsumerGroup)内的所有消费者,其订阅关系(包括过滤表达式)必须完全一致

RocketMQ Message 追加属性

1. 两种属性的本质区别

属性类型设置方法用途举例
系统属性直接setXxx()方法RocketMQ 框架内部识别和使用。例如延迟级别、Keys、Flag 等。msg.setKeys("order_123")msg.setDelayTimeLevel(3)
用户属性putUserProperty(key, value)业务自定义,存放在消息的properties字段中。用于消息过滤(Tag/SQL)或携带业务元数据。msg.putUserProperty("age", "18")msg.putUserProperty("region", "shanghai")

putUserProperty(Key, Value),Key 和 Value 都必须是 String,且 Key 禁止"__"开头。

并非所有系统属性都开放给了SQL过滤。像KEYSDELAYRETRY_TOPIC内部系统属性,通常无法直接在SQL表达式中使用

但是UserProperty都可以用来sql过滤

维度具体内容关键细节 / 避坑指南
接口全称java.io.Serializable位于java.io包,属于标记接口(内部无任何抽象方法,只起标识作用)。
接口作用允许类的对象进行序列化(对象→字节流)和反序列化(字节流→对象)。没有此标记的类,在写出或传输时会直接抛出NotSerializableException
核心常量serialVersionUID(序列化版本号)
private static final long serialVersionUID = 1L;
必须手动显式定义!若不写,类结构一变(如增删字段),JVM自动生成的哈希值就变了,反序列化必报InvalidClassException
排除字段transient关键字修饰敏感字段(如密码、Token),使其不被序列化。反序列化后该字段恢复为默认值(null/0)。
继承规则• 父类实现了Serializable→ 子类自动可序列化。
• 父类未实现 → 子类反序列化时会调用父类的无参构造器
若父类未实现且没有无参构造器,反序列化直接报错。

Springboot中快速使用RocketMQ:

生产者发送消息:

方法名大白话解释要不要等结果?有返回值吗?参数能直接传对象吗?什么时候用?
syncSend同步发送:消息发出去,原地等着,直到服务器说“收到了”才往下走。要等(阻塞)✅ 有(能拿到消息ID)❌ 不行,得手动转成字符串重要消息(比如下单、付款),必须确认发送成功才能继续。
asyncSend异步发送:消息发出去,不管了,先往下执行。等服务器回复了,再回头通过“回调”通知你。不等(非阻塞)❌ 没有(结果在回调里拿)❌ 不行,得手动转成字符串不重要但量大的消息(比如记日志),不能让它拖慢主流程。
convertAndSend转换后发送:直接把 Java 对象(比如UserOrder)扔进去,它自动转成 JSON再发出去。要等(同 syncSend)❌ 没有(直接传对象就行)懒得手动转 JSON 的时候,直接把对象往里一丢,最省事。
// 1. syncSend:等结果,手动转JSON String json = "{\"name\":\"张三\"}"; SendResult result = rocketMQTemplate.syncSend("topic", json); System.out.println("消息ID:" + result.getMsgId()); // 立刻能拿到 // 2. asyncSend:不等,手动转JSON,结果靠回调 rocketMQTemplate.asyncSend("topic", json, new SendCallback() { // 成功或失败都在这里处理,但主线程已经往下跑了 }); // 3. convertAndSend:不用转JSON,直接传对象,但没返回值 rocketMQTemplate.convertAndSend("topic", new User("张三")); // 发就完了,不用管结果

单向消息:

rocketMQTemplate.sendOneWay(topic, "这是一条单向消息");

延时消息:

在RocketMQ4.x中不支持精确的延迟时间:

/ 发送延迟消息,最后一个参数 '3' 就是延迟级别 // 级别3对应延迟10秒 (详见附录的级别表) rocketMQTemplate.syncSend(topic, MessageBuilder.withPayload(message).build(), 3000, // 发送超时时间,单位毫秒 3); // ⬅️ 延迟级别:3 代表延迟 10 秒

从5.0开始支持精确延迟时间:

// 延迟时间,单位是毫秒 // 5000 = 5秒,300000 = 5分钟,1800000 = 30分钟 long delayInMilliseconds = 300000; // 5分钟 // 使用 syncSendDelayTimeMills 发送延迟消息 rocketMQTemplate.syncSendDelayTimeMills(topic, message, delayInMilliseconds);

批量消息:

直接发送MessageList即可

消费者消费数据:

import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; // @Component:把这个类交给Spring管理(必须加) // @RocketMQMessageListener:告诉Spring我要监听哪个Topic @Component @RocketMQMessageListener( topic = "my-topic", // 监听哪个主题(必须和生产者发的一样) consumerGroup = "my-consumer-group" // 消费者组名(必须填,唯一就行) ) public class MyConsumer implements RocketMQListener<String> { // 实现这个接口的 onMessage 方法 // 只要生产者发了消息,这个方法就会被自动调用,消息内容就是参数 message @Override public void onMessage(String message) { // 写你的业务逻辑 System.out.println("收到消息啦:" + message); // 比如:保存到数据库、调用其他接口... } }

要做tag过滤@RocketMQMessageListener注解里加上selectorExpression参数

要做sql过滤RocketMQMessageListener注解加上selectorType = SelectorType.SQL92,搭配selectorExpression即可。

要把同组内消费模式改为广播在注解内加上messageModel = MessageModel.BROADCASTING即可。

如果生产者不指定MessageQueueSelector,RocketMQ 默认采用轮询(Round-robin)随机策略发送消息。

  • 订单号 A 的第1条消息(创建)发到了队列0

  • 订单号 A 的第2条消息(支付)发到了队列1
    由于队列0和队列1可能被不同的消费者实例(或不同的消费线程)并行拉取,且处理速度不同,队列0的消息可能还没处理完,队列1的支付消息已经被消费了,这就产生了乱序。

发送端:MessageQueueSelector的作用

它的作用是将同一个业务标识(如订单ID)的消息,固定路由到同一个队列

// 正确用法示例 producer.send(msg, new MessageQueueSelector() { @Override public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) { // arg 即为传入的订单ID String orderId = (String) arg; // 对订单ID取模,确保同一个订单永远进同一个队列 int index = Math.abs(orderId.hashCode()) % mqs.size(); return mqs.get(index); } }, orderId);

注意:如果队列数量发生变化(比如扩容),取模结果会变,导致同一个订单落到不同队列,因此顺序敏感的场景下,生产环境严禁动态调整队列数量

消费端:必须配合“顺序消费模式”(最关键的一环)

即使发送端用 Selector 把消息都发到了队列0,如果消费端用错模式,依然会乱序!

  • 错误做法:使用MessageListenerConcurrently(并发消费)。队列0里的消息虽然是有序的,但消费端会用多线程并发拉取,线程A处理消息1(耗时2秒),线程B处理消息2(耗时0.1秒),消息2先处理完,业务上就乱序了

  • 正确做法:必须使用MessageListenerOrderly(顺序消费)。

consumer.registerMessageListener(new MessageListenerOrderly() { @Override public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) { // 框架底层会加锁(队列锁),确保同一个队列的消息,在同一时刻**只有一个线程**在处理 // 消息1不处理完,消息2绝对无法被拉取消费 return ConsumeOrderlyStatus.SUCCESS; } });

消费失败时的“避坑”准则(极易引发乱序)

如果在顺序消费模式中,某条消息(如消息2)处理失败,千万不能返回RECONSUME_LATER(放入重试队列)!因为重试队列是独立的,会导致消息2跑到后面去,消息3、4先被消费,造成乱序。

正确做法:返回SUSPEND_CURRENT_QUEUE_A_MOMENT(挂起当前队列),让消费者稍后(默认几秒后)重新拉取整个队列头部的消息,并且阻塞后续消息的拉取,直到这条失败的消息被成功消费。

在springboot中使用:

// 1. 必须加 @Service 注解,交给 Spring 管理 // 2. 必须加 @RocketMQTransactionListener 注解,并指定生产组(和 yml 里保持一致) @Service @RocketMQTransactionListener(txProducerGroup = "tx-order-producer-group") public class OrderTransactionListener implements RocketMQLocalTransactionListener { // 模拟注入你的业务 Service(比如操作数据库的) @Autowired private OrderService orderService; /** * 方法一:执行本地事务(这是 Broker 发送半消息后,立马调用的方法) * 入参说明: * Message msg:发来的消息体 * Object arg:发送时传进来的额外参数(比如订单号) */ @Override public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 1. 把参数 arg 强转为订单号 String orderId = (String) arg; try { // 2. 执行你真正的本地业务,比如:更新订单状态为“已支付” // (这里假设 orderService.updateOrderStatus 返回 boolean) boolean success = orderService.updateOrderStatus(orderId, "PAID"); // 3. 根据业务执行结果,告诉 Broker 是提交还是回滚 if (success) { // 本地事务成功 → 提交消息(消费者可以看到了) return RocketMQLocalTransactionState.COMMIT; } else { // 本地事务失败 → 回滚消息(丢弃,消费者看不见) return RocketMQLocalTransactionState.ROLLBACK; } } catch (Exception e) { // 如果发生了异常(比如数据库超时),不知道成功还是失败 // 返回 UNKNOWN,让 Broker 过一会来回查,再决定最终状态 return RocketMQLocalTransactionState.UNKNOWN; } } /** * 方法二:事务回查(Broker 没收到提交/回滚确认时,会主动调用这个方法) * 注意:这个方法是 Broker 主动调用的,可能被调用多次,所以查询逻辑要幂等 */ @Override public RocketMQLocalTransactionState checkLocalTransaction(Message msg) { // 1. 从消息头里取出业务 ID(之前在发送时塞进去了) String orderId = (String) msg.getHeaders().get("orderId"); // 2. 去数据库查询这个订单的真实状态 Order order = orderService.getOrderById(orderId); // 3. 根据查到的结果,回复 Broker if (order != null && "PAID".equals(order.getStatus())) { // 订单已支付 → 提交消息 return RocketMQLocalTransactionState.COMMIT; } else if (order != null && "INIT".equals(order.getStatus())) { // 订单还是初始状态 → 说明本地事务没执行成功,回滚 return RocketMQLocalTransactionState.ROLLBACK; } else { // 查不到订单,或者状态异常,继续返回 UNKNOWN,让 Broker 稍后再来回查 return RocketMQLocalTransactionState.UNKNOWN; } } }

brokerID=0表示为主节点

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询