mqbroker.cmd -n 127.0.0.1:9876 autoCreateTopicEnable=true开启确保 Broker 启动时就允许自动创建主题
| 对比维度 | Message(普通消息) | MessageExt(扩展消息) |
|---|---|---|
| 谁在用 | 生产者(Producer)发送消息时使用 | 消费者(Consumer)接收消息时使用 |
| 包含什么 | 只有Topic、Tag、业务 Body(消息体) | 包含全部 Message 内容 + 系统元数据 |
| 用途 | 只负责告诉 Broker “我要发什么内容” | 告诉消费者“你收到的这条消息,在服务器里的状态是什么” |
正确的消费者管理方式
长期运行,保持单例:正确的模式是将消费者设计为长期运行的单例。在应用启动时创建并启动它,然后让它持续监听,应用关闭时才统一关闭。
优雅关闭,只做一次:
shutdown()应该只在你确定要彻底停止消息消费并释放资源时(例如应用正常下线),作为最后的清理步骤被调用一次。避免在监听器中操作:绝对不要在
MessageListener的consumeMessage()方法内部去调用consumer.shutdown()。这会导致消费者在消息处理过程中被意外销毁。
RocketMQ的消息分发遵循“组内竞争,组间复制”原则:
同组(Group相同)=竞争关系,一条消息只投递给组内一个消费者(负载均衡)。
异组(Group不同)=复制关系,每条消息会被每个组独立消费一次(业务解耦)。
同组内的“争抢”不是随机抢单条,而是“抢队列”
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过滤。像KEYS、DELAY、RETRY_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 对象(比如User、Order)扔进去,它自动转成 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表示为主节点