RocketMQ
快来分享你的内容吧~
- 01-18 21:36·Java后端MQ简介 Message Queue 消息队列,简称MQ,是一种异步通信机制 常见用途 解耦 生产者发送消息后立即返回,消费者在合适时机处理消息,减少服务之间的影响 削峰 高并发场景下暂存请求,平滑高峰流量,避免系统过载 异步 将不需要立即处理的任务放入消息队列中异步执行,减少用户请求和响应时间 消息队列的两种模式 点对点模式 一个生产者对应一个消费者 消费者主动拉取数据,消息收到后清除消息 发布查看全文加油鸭:这篇MQ知识总结太全面了!从原理到实战部署一应俱全,连内存调优细节都考虑到了,简直是手把手教学,为你点赞!341分享
2025-11-03·Java后端
Day13 RocketMQ完结
### 集群高级特性 - **DlDger文件一致性协议** > Dledger来自于开源组织OpenMessage > - 为什么需要有一致性协议? > 在一个集群中,数据写入到server中的一个节点,然后希望从集群中任意一个节点都能读到写入的数据,实现这个功能有几个核心的问题 1. 服务稳定性 各个server状态不稳定,随时可能宕机 2. 网络抖动 server之间的网络如果发生抖动,请求就有可能丢失 3. 网速问题 server之间的网络传输速度不一致,难以保证数据顺序 4. 快速响应 在server之间相互同步数据的同时,还要快速给客户端响应结果 为解决这些问题,于是提出了一致性算法 弱一致性算法:DNS系统、Gossip协议 强一致性算法:Basic-Paxos、Multi-Paxos包括Raft(Nacos-JRaft、Kafka-KRaft、RocketMQ-Dledger)、ZAB等 > - 工作流程  Log是保存在Server上的操作日志,其中的每个条目称为Entry,Entry中的操作,最终都会落到stateMachine中,Raft算法的核心就是要保证所有节点上的**Entry顺序一致** 1. 多个server基于一致性协议,会共同选举产生一个Leader,负责响应客户端请求 2. Leader通过一致性协议,将客户端的指令转发到集群所有节点上 3. 每个节点将客户端的指令以Entry形式保存到自己的Log当中,此时Entry是uncommitted状态 4. 当多数节点共同保存了Entry后,就可以执行Entry中的客户端操作,提交到State Machine中,此时Entry更新为commited状态 > Raft为每个节点设定了三种不同的角色 **Leader**:由选举产生,少数服从多数;向Follower节点发送心跳,Follower收到心跳就不会竞选Leader;响应客户端请求,集群内的所有数据变化都从Leader开始;向Follower同步操作日志 **Follower**:参与选举投票;同步来自Leader的数据;接收来自Leader/Follower的心跳,如果长时间没有接收到心跳,就自动转为Candidate **Candidate**:没有Leader时,发起投票竞选Leader **投票流程** 所有节点起始角色都为Follower,每个节点都设定了一个随机选举过期时间Election Timeout(150ms-300ms),最早过期的那个节点会优先变成Candidate角色然后向其他节点发起投票请求,每个节点在同一任期Term内只有一次投票资格,在收到其他节点的投票响应后,Candidate会重置自身的Election Timeout,继续等待其他响应,一旦收获的票数超过集群节点数的一半,该节点就成为Leader节点,开始向其他节点发送心跳,确认自己的Leader身份,其他节点收到Leader心跳后,会自动转变为Follower状态,直到Leader心跳超时或者宕机才会触发下一次选举 > - **主从节点切换的高可用集群** > RocketMQ提供的Dledger虽然确实增加了集群的高可用,但是它是把集群选举和日志同步都一起完成的,而Dledger集群下的日志会比主从集群大很多,会增加写日志的IO负担,于是RocketMQ提供了一种Controller机制,既可以使用Raft选举机制,又可以使用原生的CommitLog日志 >  具体部署方式,参见官网 https://rocketmq.apache.org/zh/docs/deploymentOperations/03autofailover。 - **BrokerContainer容器式运行机制** > 在RocketMQ4.x版本中,一个broker就是一个进程,而broker又是分主从的,压力方式不一样,主节点负责响应请求,而从节点一般就只承担冷备的作用,这种不对等的角色导致服务器资源不能充分利用。 于是在RocketMQ5.X版本中,提供了BrokerContainer模式。在一个BrokerContainer当中可以加入多个broker(master、slave、dledger),提高资源利用率,并且可以通过交叉部署来实现节点之间对等部署 >  ```bash # 修改配置文件 vim conf/container/broker-container.conf #配置端口,用于接收mqadmin命令 #listenPort=10811 #指定namesrv #namesrvAddr=worker1:9876;worker2:9876;worker3:98767 #或指定自动获取namesrv #fetchNamesrvAddrByAddressServer=false #指定要向BrokerContainer内添加的brokerConfig路径,多个config间用“:”分隔; #不指定则只启动BrokerConainer,具体broker可通过mqadmin工具添加 #brokerConfigPaths=/app/rocketmq/rocketmq-all-5.3.0-bin-release/conf/2m-2s-async/broker-b-s.properties:/app/rocketmq/rocketmq-all-5.3.0-bin-release/conf/2m-2s-async/broker-a.properties bin/mqbrokercontainer -c broker-container.conf ```  RocketMQ5.x运行架构图 ## Kafka ### Kafka服务搭建 - 单机部署 ```bash # 1. 下载 Kafka(以 3.9.1 为例) cd /opt wget https://archive.apache.org/dist/kafka/3.9.1/kafka_2.13-3.9.1.tgz # 2. 解压并创建软链接 tar -zxvf kafka_2.13-3.9.1.tgz -C /usr/local/ ln -s /usr/local/kafka_2.13-3.9.1 /usr/local/kafka # 3. 配置环境变量 cat >> /etc/profile << 'EOF' export KAFKA_HOME=/usr/local/kafka export PATH=$PATH:$KAFKA_HOME/bin EOF source /etc/profile # 生成唯一集群 ID(保存输出值) kafka-storage.sh random-uuid # 示例输出:a9sqG57YRiW34orvtakVfg vim config/kraft/server.properties # 节点 ID(单机设为 1) # node.id=1 # 监听地址(替换为你的服务器 IP) # listeners=PLAINTEXT://:9092 # advertised.listeners=PLAINTEXT://<你的IP>:9092 # 日志存储目录(建议修改为非 /tmp 路径) # log.dirs=/var/log/kafka # 角色配置(单机同时作为 broker 和 controller) # process.roles=broker,controller # controller.listener.names=CONTROLLER # controller.quorum.voters=1@<你的IP>:9093 # 安全协议映射 # listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT # inter.broker.listener.name=PLAINTEXT # 使用上一步生成的集群 ID KAFKA_CLUSTER_ID="a9sqG57YRiW34orvtakVfg" kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c /usr/local/kafka/config/kraft/server.properties # 输出提示:Formatting /data/kafka-logs with metadata.version 3.9-IV0 # 后台启动 Kafka nohup kafka-server-start.sh /usr/local/kafka/config/kraft/server.properties > /var/log/kafka/kafka.log 2>&1 & # 验证进程 ps -ef | grep kafka # 应看到:kafka.Kafka config/kraft/server.properties # 1. 创建测试主题 kafka-topics.sh --create --topic test-topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1 # 2. 查看主题列表 kafka-topics.sh --list --bootstrap-server localhost:9092 # 3. 发送消息(在新终端) kafka-console-producer.sh --bootstrap-server 60.204.184.61:9092 --topic test-topic # 输入消息内容后按 Enter # 4. 消费消息(在新终端) kafka-console-consumer.sh --bootstrap-server 60.204.184.61:9092 --topic test-topic --from-beginning ``` ## 常见MQ问题梳理 - **如何保证消息不丢失?** > 首先要明确在使用MQ的链路中,哪些场景可能丢失消息? > 1. 生产者发到Server时,有可能因为网络原因造成数据丢失 > 2. Server主从同步时,有可能因为网络造成数据丢失 > 3. 消息数据刷盘时(可能性较小),数据由page cache写入磁盘时,服务宕机,有可能造成数据丢失 > 4. Server推送到消费者时,有可能因为网络造成数据丢失  1. 生产者发送消息如何保证不丢失? 生产者发送消息确认机制 1. RocketMQ实现有三种发消息的方式 producer.sendOneWay(msg) 异步发送,不确认是否成功,有丢失消息的风险 producer.send(msg) 同步发送,等待broker确认已收到消息再进行下一步处理 producer.send(msg, new SendCallback()) 异步发送,回调确认消息是否发送成功 也可以使用RocketMQ的事物消息,确保与数据库事物修改同步 2. Kafka实现Future<RecordMetadata> future = producer.send(record) ,生产者调用future.get()方法拿到元数据,就说明已成功 3. RabbitMQ提供了Publisher Confirms机制, //添加两个回调,一个处理ack响应,一个处理nack响应 channel.addConfirmListener(ConfirmCallback ackCallback, ConfirmCallback nackCallback); 2. Broker如何保证消息不丢失?操作系统刷盘操作,没法完全保证数据安全性 1. RocketMQ 提供了一个刷盘配置flushDiskType,SYNC_FLUSH同步刷盘,发一次消息就写一次磁盘(其实底层代码是间隔10ms),ASYNC_FLUSH异步刷盘,间隔一段固定时间刷盘,性能更稳定 2. Kafka提供了参数可设置刷盘策略,当log.flush.interval.messages为1时,就是一条消息刷一次盘 ```bash flush.ms : 多长时间进行一次强制刷盘。 log.flush.interval.messages:表示当同一个Partiton的消息条数积累到这个数量时,就会申请一次刷盘操作。默认是Long.MAX。 log.flush.interval.ms:当一个消息在内存中保留的时间,达到这个数量时,就会申请一次刷盘操作。他的默认值是空。如果这个参数配置为空,则生效的是下一个参数。 log.flush.scheduler.interval.ms:检查是否有日志文件需要进行刷盘的频率。默认也是Long.MAX。 ``` 3. RabbitMQ明确说明即便是持久化队列,也不会每条数据都刷盘,基本是交由操作系统自行处理 3. Broker主从同步怎么保证消息不丢失? 1. RocketMQ - 普通集群 由于RocketMQ普通主从集群主节点挂了之后不会自动切换主节点,未同步的数据仍旧会保留在主节点上,待主节点恢复重启后,从节点会将未同步到的消息继续同步,但也正因为如此,主节点崩溃(Q5)将没有新的消息可写入 - Dledger集群 采用CP模型,将数据在主从节点之间同步,两阶段提交CommitLog确保多数节点能够同步到消息,即使切换主节点,也能极大程度的保留数据,但是极端情况下仍旧会造成数据丢失,不过可能性非常小 2. Kafka 在Kafka集群采用的是AP模型,在集群中,如果Leader Partition服务宕机,Follower会重新选举产生一个新的Leader Partition,而所有的消息都会以Leader为准,这样即使旧的Leader重启了,也是作为Follower将保存的那份数据刷掉 4. 消费者消费怎么保证消息不丢失? 消费状态确认机制 消费者处理完成后,需要给Broker一个响应,告诉对方自己已经把消息妥善处理了,如果broker没有收到响应,就会认为消息没有处理成功,会根据重试机制再次尝试消息投递,RocketMQ和Kafka根据offset进行重新投递,而RabbitMQ的Classic Queue会将消息重新入队,正常来说,这个过程丢消息的可能性极小,消费者应该考虑重复投递的幂等性问题 5. 如果MQ服务挂了,怎么保证消息不丢失? 设计缓存降级,在进行消息投递重试失败后,考虑将消息往降级缓存队列里写,继续正常执行业务,然后不断将降级缓存里的数据投递到MQ服务,等MQ服务恢复后流程就正常了 - **怎么保证消息的顺序性?** 由于队列的FIFO性质,在同一个Message Queue中,消息是天然有序的,那么要保证整体消息的有序性(一般是局部有序),就应该保障两个方面 1. Producer 将一组有序数据写入到同一个Message Queue RocketMQ和Kafka都提供了分区机制,可以让应用程序自行决定写到哪个分区(Message Queue),即能保证同一批数据的写入有序性; 而RabbitMQ则可以通过维护Exchange与Queue之间的绑定关系,将这一组局部有序的消息转发到同一个对列中,从而保证这一组有序的消息, 2. Consumer每次集中从同一个Message Queue中拿取消息 RocketMQ 消费是多线程的,但可以提供实现MessageListenerOrderly接口来对消费线程进行并发控制,保证消费有序性 Kafka 消费某一个partition时,天生就是单线程,所以也能过保证有序性 RabbitMQ 本身就是点对点的消费模式,只要保证一个队列中只有一个consumer,就能是顺序的,如果某条消息消费失败,由于它会将消费失败的数据进行重新投递,那就无法保证顺序性了 - **如何保证消息幂等性?** 1. 生产者投递过程中的幂等性 RocketMQ 在发送消息时给每条消息分配一个唯一ID,提供这个ID判断是否重复投递 Kafka 提供幂等性配置,需要打开idempotence幂等性控制(默认是打开的,但是如果其他配置有冲突,会影响幂等性配置) 生产者投递消息不会使用MQ系统自带的处理,而是在消息本身带有唯一ID或者版本号,然后由消费者自行做幂等判断 2. 消费者消费过程中的幂等性 消费时,校验唯一ID或版本号判断是否继续消费,除此之外,需要关注重试队列和死信队列中的数据,及时排查消费失败的原因 - **如何快速处理消息积压?** > 消息积压的危害 当MQ有大量消息积压时,会严重影响服务端的性能。 对于RocketMQ和Kafka来说,消息积压能力很强,短时间积压不会有太大影响,但是如果一直得不到解决,积压数据的日志文件可能丢失。 对于RabbitMQ来说,StreamQueue流式队列处理机制类似于以上两者,Classic和Quorum处理能力弱,需及时处理积压问题 > 消息积压的根本在于消费者consumer处理效率太低,所以最核心的目标就是提升consumer消费效率,消费效率上来了问题自然迎刃而解。 对于RabbitMQ,由于是点对点模式,一条消息只能由一个consumer消费,所以只需增加consumer的数量,即可解决。 对于RocketMQ和Kafka,由于同一个消费者组下的多个Consumer需要和对应topic下的MessageQueue建立对应关系,而一个MessageQueue最多只能被一个消费者consumer消费,因此,即使增加消费者,至多也只能和topic下MessageQueue的数量相同,如果此时MessageQueue的数量本身不够多的话,也就没法使用这种方法解决了。此时可以创建一个新的临时topic,配置足量的MessageQueue,然后把Consumer实例的topic指向新topic,并紧急上线一组新的topic去消费旧topic的数据往新topic里面丢(只处理这一个逻辑),这样处理能够快很多,但也是紧急使用,后续要对原有的设计进行优化。
Day12 RocketMQ的简单使用
> 今天没学多少,接前一天 - 延迟消息 - 指定延迟级别 message.setDelayTimeLevel(3);  - 指定发送时间 message.setDeliverTimeMs(System.currentTimeMillis() + 10_000L); - 批量消息 - 多条消息合并成一批发出去,减少网络IO,提升消息发送的吞吐量 - 批量消息最好不要超过1M,同一批消息的topic必须相同,且不支持延迟 - 事物消息 - 通过RocketMQ的事物机制,保证上下游数据一致性,重点是监听数据库事物,根据数据库的事物判断MQ是否投递给消费者 **示例代码请下载源码查看**org.apache.rocketmq.example.transaction.TransactionProducer  1. 生产者将消息发送至Apache RocketMQ服务端。 2. Apache RocketMQ服务端将消息持久化成功之后,向生产者返回Ack确认消息已经发送成功,此时消息被标记为"暂不能投递",这种状态下的消息即为半事务消息。 > 半消息对消费者不可见,实际上是将消息放进来一个叫RMQ_SYS_TRANS_HALF_TOPIC的系统topic > 3. 生产者开始执行本地事务逻辑。 4. 生产者根据本地事务执行结果向服务端提交二次确认结果(Commit或是Rollback),服务端收到确认结果后处理逻辑如下: - 二次确认结果为Commit:服务端将半事务消息标记为可投递,并投递给消费者。 - 二次确认结果为Rollback:服务端将回滚事务,不会将半事务消息投递给消费者。 5. 在断网或者是生产者应用重启的特殊情况下,若服务端未收到发送者提交的二次确认结果,或服务端收到的二次确认结果为Unknown未知状态,经过固定时间后,服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查。 > 回查次数通过transactionCheckMax参数设置,默认15次,回查间隔通过transactionCheckInterval参数设置,默认60s > 6. 生产者收到消息回查后,需要检查对应消息的本地事务执行的最终结果。 7. 生产者根据检查到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤4对半事务消息进行处理。 - ACL权限机制 - topic权限 通过perm字段配置 2:禁写禁订阅,4:可订阅,不能写,6:可写可订阅 - broker权限配置 1. 在broker.conf 设置aclEnable=true 2. 修改plain_acl.yml(热加载,不需要重启)  3. 客户端使用,引入org.apache.rocketmq.rocketmq-acl包,声明时传入AclClientRPCHook对象即可 - SpringBoot整合RocketMQ使用 ```xml <dependency> <groupId>org.apache.rocketmq</groupId> <artifactId>rocketmq-spring-boot-starter</artifactId> <version>2.3.0</version> <!-- 推荐使用最新稳定版、注意版本兼容性 --> </dependency> ``` ```yml rocketmq: name-server: localhost:9876 # 必填!NameServer 地址 producer: group: my-producer-group # Producer 组名(必须唯一) send-message-timeout: 3000 # 发送超时(毫秒) retry-times-when-send-failed: 2 # 同步发送失败重试次数 # consumer: # 如果需要消费消息,再配置 consumer # listeners: # your-topic: your-consumer-group ``` ```java // Producer import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @Service public class MessageService { @Autowired private RocketMQTemplate rocketMQTemplate; public void sendMessage(String topic, String message) { // 同步发送:topic + 消息体(自动序列化为 JSON) rocketMQTemplate.convertAndSend(topic, message); // 也可以带 tag:topic:tag // rocketMQTemplate.convertAndSend("OrderTopic:pay", order); } } // Consumer import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; /** * 参考官方测试类 com.roy.rocketmq.SpringRocketTest * 一个RocketMQTemplate实例只能包含一个生产者,也只能往一个topic下发送消息,如果需要往另一个topic下发消息 * 就要通过@ExtRocketMQTemplateConfiguration()声明一个子实例 * 事物消息使用 @RocketMQTransactionListener,通过rocketMQTemplateBeanName指向具体子类 */ @Component @RocketMQMessageListener( topic = "TestTopic", consumerGroup = "my-consumer-group" // 必须与 producer 的 group 不同 ) public class TestTopicConsumer implements RocketMQListener<String> { @Override public void onMessage(String message) { System.out.println("Received message: " + message); // 处理业务逻辑 } } ```
Day11 RocketMQ集群部署客户端使用
- 集群部署(由于服务器资源限制2G,选择部署三主) > 我搭建步骤是对的,但是因为太吃内存了,2G支撑不起,所以直接game over > - 2m-noslave 两主无从节点,存在单点故障【每台机器开一个主节点,勉强能跑】 - 2m-2s-sync/async 两主两从同步/异步,主节点宕机需手动切换【在每台机器上要开一主一丛,我2G内存受限,会OOM】 - dledger 具备主从切换的高可用集群。集群中的节点会基于Raft协议随机选举出一个leader【比两主两从还更重,至少需要4G内存】 - 基于Raft协议,确保多个副本之间特别数据,保证数据强一致性 - 会接管RocketMQ原生的文件写入,写入速度稍慢,会有性能影响 ```bash # 将单机部署的机器停机 kill -9 pid # 在机器 A 上 cp /opt/rocketmq/conf/broker.conf /opt/rocketmq/conf/broker-tmp.properties vim /opt/rocketmq/conf/broker-template.properties #brokerClusterName=MyRocketMQCluster #brokerId=0 #deleteWhen=04 #fileReservedTime=48 #brokerRole=ASYNC_MASTER #flushDiskType=ASYNC_FLUSH #listenPort=10911 #storePathRootDir=/opt/rocketmq/store #storePathCommitLog=/opt/rocketmq/store/commitlog # namesrvAddr 先留空,后面统一替换 # 打包mq tar -czf rocketmq.tar.gz rocketmq # 分发到BC两台机器 # 分发到机器 B scp rocketmq.tar.gz root@192.168.1.102:/opt/ # 分发到机器 C scp rocketmq.tar.gz root@192.168.1.103:/opt/ # 解压并删除原文件 cd /opt tar -xzvf rocketmq.tar.gz rm -f rocketmq.tar.gz cd rocketmq mkdir -p store/{commitlog,consumequeue,index,checkpoint,abort} mkdir -p /var/log/rabbitmq # 复制编辑配置文件,三台机器都要改 cp /opt/rocketmq/conf/broker-tmp.properties /opt/rocketmq/conf/broker-c.properties vim /opt/rocketmq/conf/broker-c.properties # brokerName=broker-b # namesrvAddr= IP1:9876;IP2:9876;IP3:9876 # 看内存情况,禁用以下配置,机器4G以上且只部署MQ服务无需操作 # enableScheduleMessageStats=false # timerWheelEnable=false # timerColdDataCheckEnable=false # enableLmq=false # enableMultiDispatch=false # transientStorePoolEnable=false # sendMessageThreadPoolNums=1 # pullMessageThreadPoolNums=1 # 把临时文件删除 rm -f /opt/rocketmq/conf/broker-tmp.properties # 启动nameserver nohup sh /opt/rocketmq/bin/mqnamesrv > /var/log/rocketmq/namesrv.log 2>&1 & # 分别启动broker nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-a.properties > /var/log/rocketmq/broker.log 2>&1 & nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-b.properties > /var/log/rocketmq/broker.log 2>&1 & nohup /opt/rocketmq/bin/mqbroker -c /opt/rocketmq/conf/broker-c.properties > /var/log/rocketmq/broker.log 2>&1 & # PS:不出意外的话得出意外了,我2G的内存不够,启动不了,即使禁用了一些配置,AI推荐降级到4.9,但我已经乏了 # 验证集群 # 查看集群 /opt/rocketmq/bin/mqadmin clusterList -n 192.168.1.101:9876 # 应输出: # MyRocketMQCluster broker-a 192.168.1.101:10911 # MyRocketMQCluster broker-b 192.168.1.102:10911 # MyRocketMQCluster broker-c 192.168.1.103:10911 # 测试发送 export NAMESRV_ADDR=192.168.1.101:9876 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer ```     - 运行时架构图(核心:nameServer、broker、client) nameServer:服务注册协调中心,独立启动 broker:提供消息存储、传递、查询等功能,是RocketMQ中最繁琐也是最娇贵的组建 client:客户端 > 💡为什么RocketMQ不使用Zookeeper作为注册中心,而是自己实现nameServer? > 1. 设计简洁且专用 专门的场景设计,更加轻量,易部署 > 2. 高性能 zookeeper比较重,繁杂的管理机制和强一致性带来了较多开销 > 3. 高可用 无状态,多个nameserver之间对等,即使只有一台示例也不会影响整体运行,而zookpeer在节点之间同步要求严格 > 4. 降低依赖性 自己实现可以降低外部系统的依赖,简化系统复杂度,同时能够自己掌控优化 > 5. 定制化需求 需要实现特定的功能,不必受限与zookeeper的实现和接口  ### 客户端编程模型 - 消息处理模型  - 消息确认机制 > 要支持互联网金融场景,消息安全必须是最高优先级保障。 消息安全有两方面要求,一方面是生产者要能确保将消息发送到Broker上,另一方面消费者要能确保从broker上获取到消息 > - 消息生产端采用消息确认加重试机制保障消息正常发送到RocketMQ > 三种发送消息的方式 1. **单向发送SendOneWay** 只管发送消息,不管成不成功。发送消息效率高,但如果发送失败则无法补救,适用于一些追求效率且允许消息丢失的业务场景。 2. **同步发送send** 发送完消息要等待broker回复结果{SEND_OK,FLUSH_DISK_TIMEOUT,FLUSH_SLAVE_TIMEOUT,SLAVE_NOT_AVAILABLE},保证消息一定发到broker,如果返回失败,可以进行重试,但返回失败不一定代表没有推送给下游消费者。能保障消息发送的安全性,但效率低,会阻塞当前线程。 3. 异步发送send callback{onSuccess,onException} 不阻塞主线程,由异步回调处理成功或失败的情况,但是在回调结束之前,生产者主线程不能调用shutDown方法关闭主线程。 > - 消费者端采用状态确认机制保证消费者一定能正常处理对应消息 > CONSUME_SUCCESS | RECONSUME_LATER 当broker拿到的是RECONSUME_LATER状态时,会将当前消息放到对应消费者组的重试topic中,以免影响正常队列的运行,当重试次数大于16时,则将消息推入到死信topic,如仍需要重试,可以人工介入进行消息补救 > - 消费者组可以自行指定起始消费位点 > CONSUME_FROM_LAST_OFFSET | CONSUME_FROM_FIRST_OFFSET | CONSUME_FROM_TIMESTAMP 虽然提供了这么一种功能,但是消费者无法知道自己是在那里中断消费的,所以最好是建一个新的消费者组,再从起始点开始消费,也可以指定具有的时间点,consumer.setConsumerTimestamp("20251223171201"); > - 消息处理模式 - 广播模式 一个消息,推送到所有消费者实例,不关心消费者组 - `consumer.setMessageModel(MessageModel.BROADCASTING);` - 将offset交由consumer自行保管,只要consumer拉取,就返回对应消息,但是broker端不会对消费失败的消息进行重试 - 集群模式(默认) 一个消息只会推送到指定的消费者组,由消费者组中的多个实例共同消费 - 给每个ConsumerGroup维护一个统一的offset,同一个组内只会被消费一次 - 消息过滤 - 简单过滤 订阅时使用tag过滤自己感兴趣的内容 - consumer.subscribe("TagFilterTest", "Tag"); - 使用两个竖线(||)连接多个Tag值,也可以使用星号(*)匹配所有 - SQL过滤 使用标准SQL语句进行过滤 WHERE 后面的过滤条件 - 数值比较,比如:**>,>=,<,<=,BETWEEN…AND,=;** - 字符比较,比如:**=,<>,IN;** - **IS NULL** 或者 **IS NOT NULL;** - 逻辑符号 **AND,OR,NOT;** - 顺序消息机制 - 生产者将一组有序的消息发送至同一个MessageQueue,天然保证有序性 - 消费者实现MessageListenerOrderly接口,按序拿完这一批数据后才去读其他MessageQueue的消息 - 注意事项 - 有序性一般是指局部有序,如果需要全局有序,那就只能使用一个queue,性能就很差了 - 生产者尽可能将有序消息分散到不同的queue,避免数据过于集中导致热点竞争 - 消费者端只进行有限次数的重试,如果消息处理失败,要保持有序性,就得阻塞一直重试失败消息的消费,但如果一直失败,此时就无法消费其他消息了。 - 消费者如果处理出现异常,建议返回ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT挂起替代抛出异常
Day10 MQ产品选择与Rocket MQ搭建
## MQ简介 Message Queue 消息队列,简称MQ,是一种异步通信机制 ### 常见用途 - 解耦 - 生产者发送消息后立即返回,消费者在合适时机处理消息,减少服务之间的影响 - 削峰 - 高并发场景下暂存请求,平滑高峰流量,避免系统过载 - 异步 - 将不需要立即处理的任务放入消息队列中异步执行,减少用户请求和响应时间 ### 消息队列的两种模式 - 点对点模式 - 一个生产者对应一个消费者 - 消费者主动拉取数据,消息收到后清除消息 - 发布订阅模式 - 可以有多个topic - 消费者消费数据后不删除数据 - 每个消费者相互独立,都可以消费到数据 ### 常见MQ产品对比 > 大吞吐量大数据优先选Kafka,灵活低延迟选RabbitMQ,可靠事物国产化选RocketMQ > | 对比维度 | RabbitMQ | Kafka | RocketMQ | | --- | --- | --- | --- | | 开源背景 | Pivotal(VMware) | LinkedIn → Apache | 阿里巴巴 → Apache | | 开发语言 | Erlang | Scala + Java | Java | | 核心模型 | Exchange + Queue(灵活路由) | Topic + Partition(日志流) | Topic + MessageQueue(轻量注册中心) | | 协议支持 | AMQP 0.9.1/1.0, MQTT, STOMP | 自定义二进制协议 | 自定义协议(部分兼容 OpenMessaging) | | 吞吐能力 | 1–3 万 QPS | 50 万+ QPS(高吞吐) | 10–50 万 QPS(可调优至百万级) | | 端到端延迟 | 50μs – 1ms(低延迟) | 5–20ms(批量优化后) | 1–10ms | | 消息持久化 | 支持(需显式开启) | 天然持久(顺序写磁盘) | 天然持久(CommitLog 顺序写) | | 复制机制 | 镜像队列(主从同步) | ISR 副本(In-Sync Replicas) | 主从同步(支持同步/异步刷盘 + 复制) | | 一致性保证 | 强一致(镜像队列) | 最终一致(ACK=all 可近似强一致) | 强一致(金融级设计) | | 可靠性(默认) | 中 | 中(依赖配置) | 高 | | 顺序消息 | 单队列内有序 | 分区内有序 | 全局有序 / 分区有序 | | 事务消息 | AMQP 事务(性能差) | 生产者事务(仅原子性) | 完整二阶段事务消息 | | 延迟消息 | 插件支持(18 级固定延迟) | 不支持(需外部实现) | 原生支持(18 级 + 精确延迟) | | 死信队列 | 支持 | (通过重试 Topic 模拟) | 支持 | | 消息过滤 | Exchange 路由 | 消费组级别 | Tag + SQL92 表达式 | | 消息轨迹 | 需插件) | 部分(依赖第三方) | 原生支持 | | 多语言客户端 | 极好(官方支持全) | 极好(社区丰富) | 偏 Java(其他语言生态较弱) | | 管理界面 | 内置 Web UI(优秀) | (依赖 Kafka Manager/UI for Kafka 等) | RocketMQ Console(功能完善) | | 大数据生态集成 | 弱 | 极强(Flink/Spark/ELK 等) | 中(阿里系生态强) | | 动态扩容 | 困难(集群结构固定) | 容易(加 Broker/Partition) | 较容易(加 Broker) | | 典型适用场景 | 企业内部系统、任务队列、IoT、微服务通信 | 日志采集、流处理、大数据管道、事件溯源 | 电商交易、金融支付、订单系统、国产化项目 | | 主要缺点 | 吞吐低、Erlang 扩展难、集群扩容不灵活 | 延迟高、无原生延迟/事务、运维复杂 | 多语言支持弱、社区生态略逊于 Kafka | ## RocktMQ 阿里巴巴开源的使用java开发的消息中间件 ### 单机服务搭建 ```bash # 下载官方二进制包 wget https://archive.apache.org/dist/rocketmq/5.1.4/rocketmq-all-5.1.4-bin-release.zip # 如果没有的话需要下载unzip,有的话跳过这一步 yum -y install unzip # 解压(必须用 unzip) unzip rocketmq-all-5.1.4-bin-release.zip mv rocketmq-all-5.1.4-bin-release /opt/rocketmq # 由于我的服务器内存只有2G,所以必须把内存改小 # 修改NameServer内存 vim /opt/rocketmq/bin/runserver.sh # 将 # JAVA_OPT="${JAVA_OPT} -server -Xms4g -Xmx4g -Xmn2g -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m" # 改为 # JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m -XX:MetaspaceSize=64m -XX:MaxMetaspaceSize=128m" # 开发测试自己玩够用 # 修改Broker内存 vim /opt/rocketmq/bin/runbroker.sh # JAVA_OPT="${JAVA_OPT} -server -Xms8g -Xmx8g" # 改为 # JAVA_OPT="${JAVA_OPT} -server -Xms512m -Xmx512m" # 修改jvm启动内存 vim /opt/rocketmq/bin/tools.sh # JAVA_OPT="${JAVA_OPT} -server -Xms1g -Xmx1g -Xmn256m -XX:MetaspaceSize=128m -XX:MaxMetaspaceSize=320m" # 改为 # JAVA_OPT="${JAVA_OPT} -server -Xms128m -Xmx128m -XX:MetaspaceSize=64m -XX:MaxMetaspaceSize=128m" # 修改环境变量 export NAMESRV_ADDR='localhost:9876' # 启动 NameServer(默认端口 9876) nohup sh /opt/rocketmq/bin/mqnamesrv > /var/log/rocketmq/namesrv.log 2>&1 & # 启动 Broker(连接本地 NameServer) nohup sh /opt/rocketmq/bin/mqbroker -n localhost:9876 > /var/log/rocketmq/broker.log 2>&1 & # 查看是否启动成功 jps # BrokerStartup # NamesrvStartup # 测试发送消息 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Producer # 测试接收消息 sh /opt/rocketmq/bin/tools.sh org.apache.rocketmq.example.quickstart.Consumer # 消息能被消费,无报错即正常 ``` 
消息队列 - 理论梳理
本文系统介绍消息队列概念,以及 RabbitMQ , RocketMQ , kafka 三个消息队列的核心概念。 ## 为什么需要消息队列? 在单体项目中,如果没有消息队列,那么: 1. 上游(用户)的操作(如生成视频)需要系统响应很久,那么上游就会长时间等待响应,不能做别的操作,阻塞线程。 2. 大量用户使用生成视频功能,每个任务下游直接处理(因为中间没有任何缓冲),可能导致下游服务处理不过来崩溃。 3. 系统可能在高峰时间段流量会突增,大部分时间处理的过来,如果盲目加机器,会造成成本效益低。 在分布式系统中,不同服务或模块之间需要通信。同步调用的缺点如下: 1. **系统耦合度过高**:服务间直接依赖,就像用胶水粘在一起。一个服务的变更或故障,可能直接影响到其他服务,维护和扩展变得困难。 2. **同步阻塞导致性能瓶颈**:主流程需要等待所有依赖操作完成才能继续。比如用户注册后,同步发送邮件、短信、初始化积分等,每一步都可能耗时,导致用户响应时间很长。 3. **无法应对突发流量(峰值冲击)** :在秒杀、大促等场景,瞬时流量远超系统处理能力。所有请求直接压到数据库等核心服务,极易导致系统崩溃。 **消息队列(MQ)** 的出现,就是为了解决这些痛点。它像是一个“中间人”或“缓冲带”,让服务间的通信更灵活、更可靠。 一般要用到消息队列多半是在分布式系统下。 ## 消息队列 ### 介绍 你可以把消息队列想象成一个**智能的邮局或消息中转站,临时(也可以持久)帮你存储待处理的任务。** - **消息**:便是待执行的任务,比如生成视频任务太耗时了,可以先放进消息队列。**消息本质上是一个自定义的 DTO对象**、JSON数据,并不是真的传递一个Task到消息队列,传递的是任务的描述信息,消费者接受后会根据信息执行任务。 - **生产者**:产生任务的上游,可以是用户,上游服务等等,生产者是一个抽象的概念,任何往消息队列放任务的角色都是生产者。 - **消费者**:处理消息队列内任务的服务,同样是抽象概念。 - **队列**:负责存储任务的容器,可以是 redis数据库,也可以用集合来手搓一个容器,最常用的是 rabbitMQ, RocketMQ , kafka  消息队列有两个含义: **指代整个技术/系统:**当我们说:“我们项目里用了消息队列”,这里指的是一整套异步通信的技术和解决方案。 **指代具体的数据结构:**当我们说:“消息被发送到消息队列里”,这里特指那个存储消息的 FIFO 数据结构,也就是**队列**。 ### 消息队列的作用 消息队列的核心价值主要体现在五个方面 **解耦**:生产者和消费者无需直接知道对方的存在,只约定好消息格式即可。新增或减少一个消费者,生产者完全不用改代码,就像我们不用知道快递小哥的存在,我们只要知道菜鸟驿站在哪就行。 **异步**:生产者发送消息后立即返回,不用等待消费者处理。非核心操作(如发送邮件、记录日志)异步处理,大大提升主流程响应速度和用户体验 **削峰填谷**:在流量高峰期,消息队列充当缓冲区,暂存大量请求。后端服务按照自己的处理能力平稳地从队列中消费消息,避免被瞬间洪流冲垮 **可靠性:**消息队列一般会有持久化消息和防止消息丢失的功能。 **顺序性:**保证消息按顺序消费,主流消息队列都有顺序消费的实现。 ### 消息队列的模型 点对点(网上也说集群模式):一条消息只能有一个消费者消费,这条消息不能重复消费。  发布/订阅( 广播模型 ) : 消息队列会把一条消息广播给所有订阅了自己的消费者,都要老老实实处理。  ### 消息是什么? 我们知道生成者和消费者之间通过消息通信,那么消息是什么?是什么数据类型? 当我们说发送一个任务给消费者(服务端)执行的时候,发送的其实不是任务本身,而是一段任务描述或者执行任务所需的数据,服务端接受到任务的描述或者数据,就会在本地任务中执行。注意消息本身不是 Task 任务, 而是数据容器,某个数据类型。 无论你要发送什么数据类型,底层都会序列化成JSON 并转成 byte[] 数组传输,凡是能序列化JSON 或者 能转成 byte[] 的数据类型,都可以发送。 ```java Producer 任何数据类型 → JSON → byte[] → MQ → byte[] → Consumer JSON → 可选反序列化为Java对象 ``` 我们来看RocketMQ 能生产什么消息,看看如何发送。 String 类型消息: ```java public static void main(String[] args) throws Exception { // 声明一个默认的生成者 DefaultMQProducer producer = new DefaultMQProducer("example-producer-group"); // 绑定看板 producer.setNamesrvAddr("127.0.0.1:9876"); producer.start(); // 声明一条 String 类型消息 String strMsg = "Hello RocketMQ"; // 消息转成 byte[] 了 Message msg1 = new Message("TopicTest", "TagA", strMsg.getBytes("UTF-8")); producer.send(msg1); producer.shutdown(); } ``` 自定义对象转成 JSON ,作为一条信息: ```java class OrderDTO { String orderId; int count; } OrderDTO dto = new OrderDTO("A001", 10); byte[] body = JSON.toJSONString(dto).getBytes(StandardCharsets.UTF_8); // byte[] body = JSON.toJSONBytes(dto); 这样也行 producer.send(body); ``` Map / List /Set 等java集合: ```java Map<String, Object> map = new HashMap<>(); map.put("a", 1); map.put("b", "Hello"); byte[] body = JSON.toJSONBytes(map); producer.send(body); ``` 发送 byte[] 本身: ```java byte[] body = new byte[]{1,2,3,4}; producer.send(body); byte[] body = Files.readAllBytes(Paths.get("test.png")); producer.send(body); ``` 发送 XML : ```java String xml = "<user><id>1</id><name>Tom</name></user>"; byte[] body = xml.getBytes(StandardCharsets.UTF_8); producer.send(body); ``` 发送 CSV : ```java String csv = "id,name\n1,Tom"; byte[] body = csv.getBytes(StandardCharsets.UTF_8); producer.send(body); ``` 发送加密内容: ```java byte[] plain = "secret".getBytes(StandardCharsets.UTF_8); byte[] encrypted = encrypt(plain); producer.send(encrypted); ``` 发送压缩数据(gzip): ```java byte[] original = "Hello MQ".getBytes(StandardCharsets.UTF_8); ByteArrayOutputStream bos = new ByteArrayOutputStream(); GZIPOutputStream gzip = new GZIPOutputStream(bos); gzip.write(original); gzip.close(); byte[] body = bos.toByteArray(); producer.send(body); ``` 发送 Kryo 序列化后的对象: kryo 是一个序列化库,将java对象转换为可以网络传输的对象,不用Kryo也可以用Java原生序列化,别的序列化库。 ```java Kryo kryo = new Kryo(); // 创建 Kryo 实例 ByteArrayOutputStream bo = new ByteArrayOutputStream(); Output output = new Output(bo); // 绑定输出流 kryo.writeObject(output, dto); // dto 对象序列化成二进制,写入 output output.close(); // 关闭流 byte[] body = bo.toByteArray(); // 得到最终 byte[] mq.send(body); // 发送到消息队列 ``` **可以发送的消息包括:** - String - JSON - DTO / POJO - Map / List - 原始 byte[] - 文件(图片、视频、PDF、zip) - XML - CSV - Kryo 序列化对象 - Java 原生序列化对象 - 加密后的数据 - 压缩数据 本质上是参数的跨服务传递。 ### Rabbit MQ 消息队列概念汇总 RabbitMQ 和 RocketMQ 关键概念区别:RabbitMQ没有 **NameServer** **,** 使用exchange + binding key + routing key 实现路由。 RabbitMQ 的集群信息和 Broker 数量由集群内部自己维护,内部是Erlang写的,Erlang自带一个轻量级数据库,维护节点元信息,所以不需要 NameServer 。 **生产者**:生产消息给队列,生产者绑定交换机,只能往交换机生产消息,生产者发布消息只能指定一个交换机,路由键可以是通配符。 **交换机**(**exchange** ):交换机通过路由键和绑定键的匹配,将生产者的发来的消息路由到队列中。有多种类型交换机。  **消费者**:消费队列的消息,没有消费者组的概念,消费者只能绑定队列,一个消费者可以监听多个队列。 **队列(queue)**:存储消息的容器,Broker 管理的就是队列,通过绑定建,将队列绑定到交换机,可以声明绑定多个交换机。 **routing key** : 路由键,在消息形成时指定路由键,路由键可以是通配符,一条消息发给多个队列。 **binding key** : 绑定键, 队列在绑定到 Exchange 时所设置的匹配规则, Exchange 根据 routing key 与 binding key 是否匹配,决定是否将消息投递给该队列。可以是通配符 **通配符规则:让消息的传递变的很灵活。** | 通配符 | 含义 | 示例 | | ------ | ----------------------- | ------------------------------------------------------------ | | `*` | 匹配 **恰好一个单词** | 绑定键:`china.*.weather` 匹配:`china.news.weather` ✅ 不匹配:`china.weather` ❌ | | `#` | 匹配 **零个或多个单词** | 绑定键:`china.#` 匹配:`china.news` ✅、`china.news.weather` ✅、`china` ✅ | **Connection** : TCP 连接,由客户端(Producer 或 Consumer)建立 ,一个 TCP 连接是单个通道,消息拆成多个包,那只能一个包一个包发送,而且只能先把这个消息的包发完,再发送下一个消息的包, 如果客户端多个线程同时往这个 TCP 发送消息,数据包会互相混在一起。如果开通多个Connection连接, 会多次 TCP 握手,消耗资源,于是 Channel 就来拯救这种情况。 **Channel(通道)**:在一个连接里可以开多个轻量级通道,通道负责发送/接收消息。 这是RabbiMQ做的优化,一个逻辑概念,不是物理通道,给消息绑定一个 Channel_id 和channel会话实现,这样一个TCP连接 + 多个Channel 实现了多线程同时往 TCP 写消息。 **channel会话**:一个逻辑概念,本质上是一个数据结构实现,一个会话对象存储了 : - **Channel 状态**:open / close - **事务状态**(transaction) - **消息确认状态**(哪些消息已经 ack/nack) - **QoS(prefetch)**:当前可以未 ack 的消息数 - **绑定的队列、交换机信息** 把会话对象和Channel_id 绑定在消息中,便实现了一个 TCP + 多个Channel , 这样多条消息的数据包可以混合、交替发送, 最大化利用单个 TCP 连接资源 ,榨干 TCP 。 #### RabbitMQ 属性(队列类型) rabbitMQ 的队列可以设置属性,让队列具备某中特性,属性可以组合搭配: **1)持久化队列(durable=true) :** Broker 重启后队列仍存在, Broker 在内存中维护队列对象,同时写入 **磁盘文件。** - 原理:消息设置 delivery_mode=2 ,便也支持持久化, Broker 收到消息后 , 将消息写入磁盘 , 仅当消息安全写入磁盘后才返回 ack 给 Producer 。 - 当 **消费者消费并确认(ack)消息** 后,Broker 才会把消息从队列和磁盘上删除。 使用场景: 关键业务消息 (电商订单,物流消息推送) , 任务队列 (视频生成,图片处理等人物)都需要可靠性。 **2)独占队列(exclusive=true):** 一种 只能被创建它的 Connection 使用的队列。 - 原理:当客户端(Connection)声明一个 **独占队列,** Broker 会在内存中为这个队列分配一个逻辑对象,也就是connection对象的引用,以此来维护队列和连接的关联性。 - - 队列一旦创建,就和创建它的 Connection 绑定。 - 当这个 Connection 关闭或断开时,队列会被自动删除(如果同时设置了 `auto-delete`)。 - 只有创建该队列的 Connection 可以声明、消费消息。 - 其他 Connection 无法访问(无法订阅或发送消息到这个队列)。 使用场景: 临时队列场景,在线客服系统 (每次打开客户咨询是新的对话), 多人协作软件、游戏实时状态更新 **3)自动删除队列(autoDelete=true):** 当 最后一个消费者断开后,队列会被 Broker 自动删除。 - 原理: Broker 为队列维护一个 **消费者计数,**当消费者全部断开, Broker 检查 `autoDelete=true`, 队列被自动删除 使用场景:和独占队列一样适用于临时队列,不同的是队列销毁方式。 **4)临时队列:**不指定队列名, Broker 自动生成一个唯一队列名 , 通常格式类似 `amq.gen-<随机字符串>` - 原理: 队列对象存储在 Broker 内存中,绑定以下属性: - - `exclusive=true` → 队列绑定创建它的 Connection - `autoDelete=true` → Connection 关闭时自动删除 使用场景:这是 匿名临时独占队列 ,临时队列使用场景都相似,多一种技术选型。 #### 消息特性 决定可靠性、消费顺序和优先级 | 概念 | 定义 | 作用 / 使用场景 | | ------------------------------- | ------------------------------------------------------------ | ------------------------------------------- | | **持久化(Delivery Mode=2)** | 消息写入磁盘,保证 Broker 重启后不会丢失 | 关键业务消息,如订单、支付、任务队列 | | **非持久化(Delivery Mode=1)** | 消息只存在内存,Broker 重启会丢失 | 临时通知、日志、实时数据 | | **消息确认(Ack / Nack)** | 消费者处理消息后向 Broker 确认,Broker 才删除消息 | 确保消息被正确消费;支持 at-least-once 语义 | | **预取 / QoS(prefetch)** | 每个消费者一次可以接收的未 ack 消息数量 | 控制消费速率,防止消费者处理不过来 | | **TTL(Time-To-Live)** | 消息过期时间,到期后自动删除 | 临时消息、延迟消息 | | **优先级消息** | 消息带优先级,消费者先消费高优先级消息,**只能发送给优先队列** | 异步任务队列、紧急事件处理 | | **死信(Dead Letter)** | 无法正常消费的消息(拒绝、过期、队列满)被转入 DLX | 记录异常消息,做后续补偿或监控 | #### 交换机(Exchange) 决定消息路由策略 | 类型 | 定义 | 消息路由方式 | 使用场景 | | ------------------------------------ | ------------------------------------------------ | -------------------------------------- | -------------------------- | | **Direct Exchange** | 直连交换机 | 根据 `routing key` 精确匹配队列 | 单点消息投递,例如任务队列 | | **Fanout Exchange** | 扇出交换机 | 广播消息到绑定的所有队列 | 广播通知、消息广播系统 | | **Topic Exchange** | 主题交换机 | 支持通配符 `*`、`#` 匹配 `routing key` | 日志系统、主题订阅 | | **Headers Exchange** | 头交换机 | 根据消息头字段匹配队列 | 灵活路由、条件匹配 | | **默认交换机(Default / nameless)** | 每个队列都绑定到默认交换机,`routing key=队列名` | 自动直连 | 简单单队列发送 | #### 消费模式 决定消息接收和确认策略 | 模式 | 定义 | 特点 / 使用场景 | | ----------------------------- | -------------------------------- | ---------------------------- | | **Push 模式** | Broker 主动推送消息给消费者 | 实时性好,常用模式 | | **Pull 模式** | 消费者主动拉取消息 | 控制消费节奏,可结合批量拉取 | | **自动 ack(auto-ack=true)** | 消费者接收消息后自动确认 | 简单快速,但存在消息丢失风险 | | **手动 ack(manual ack)** | 消费者处理完再发送 ack | 消费可靠性高,支持异常重试 | | **事务模式(tx)** | 发送或确认消息可回滚 | 确保发送原子性,但性能低 | | **Confirm 模式** | 生产者确认消息是否被 Broker 收到 | 高性能可靠发送,替代事务 | 如果把 RabbitMQ 的特性比作一个工具箱,事务就像是那个被放在角落里、布满灰尘、很少被想起来的旧工具, RocketMQ 才是使用事务的消息对立,事务能力很强。 #### 高级特性 支撑业务扩展和可靠性保障 | 特性 | 定义 | 使用场景 | | ------------------------ | -------------------------------------------------- | ------------------------ | | **死信队列(DLQ)** | 消息因拒绝、过期或队列满无法消费,进入指定死信队列 | 异常消息处理、补偿机制 | | **TTL(消息/队列过期)** | 消息或队列超过指定时间自动删除 | 延迟消息、临时消息 | | **优先级队列** | 队列支持消息优先级,高优先级消息先消费 | 紧急任务处理、异步调度 | | **镜像队列(HA Queue)** | 队列在集群节点间复制多份,保证高可用 | 集群高可用、故障恢复 | | **延迟队列 / 延时消息** | 消息延迟指定时间后再投递 | 定时任务、定时提醒 | | **批量确认 / 批量发送** | 对消息进行批量 ack 或批量发布 | 提高吞吐量,减少网络开销 | #### 生产者工作原理 生产者的核心任务是:**将消息安全、高效地发送到指定的 Exchange**。它不关心消息的去向,只关心发送动作本身。 1. **建立连接与通道**: - - 生产者首先与 RabbitMQ Broker 建立一个 **TCP 连接**。 - 在这个连接之上,它会创建一个或多个**通道**。通道是轻量级的虚拟连接,所有的操作都在通道中进行,避免了为每个操作都建立 TCP 连接的开销。 1. **开启可靠性模式(关键)**: - - 为了确保消息不丢失,生产者需要开启可靠性机制。有两种选择: - - - **事务模式**:通过 `**channel.txSelect()**` 开启。发送一批消息后,调用 `**channel.txCommit()**` 提交或 `**channel.txRollback()**` 回滚。**此模式性能极差,已不推荐使用**。 - **发布者确认模式**:通过 `**channel.confirmSelect()**` 开启。这是**异步**的高性能模式。Broker 在成功接收消息后,会异步地向生产者发送一个确认(ACK)。 1. **发布消息**: - - 生产者调用 `**channel.basicPublish()**` 方法发送消息。 - 此方法需要指定核心参数:**Exchange 名称**、**Routing Key**(路由键)和**消息体**(包含消息内容和各种属性)。 1. **接收确认**: - - 在发布者确认模式下,生产者无需阻塞等待。它可以通过监听器异步接收 Broker 返回的 `**basic.ack**`。 - 如果收到 ACK,表示消息已成功到达 Broker。如果未收到(或收到 `**nack**`),生产者可以选择重发。  #### Exchange 工作原理 Exchange 是 RabbitMQ 的**消息路由核心**。它接收来自生产者的消息,并根据**路由规则**将消息投递到一个或多个队列中。它本身不存储消息。 1. **接收消息**:Exchange 接收生产者发来的消息,并解析出其中的 **Routing Key**。 2. **匹配绑定**:Exchange 会查看与自己绑定的所有队列,以及每个队列的 **Binding Key**(绑定键)。 3. **执行路由**:根据自身的**类型**和**路由算法**,决定将消息投递到哪些队列: - - **Direct Exchange**:将 Routing Key 与 Binding Key 进行**精确匹配**。完全匹配则投递。 - **Fanout Exchange**:忽略 Routing Key,将消息广播到所有绑定的队列。 - **Topic Exchange**:将 Routing Key 与 Binding Key 进行**模式匹配**。支持 `*****`(匹配一个单词)和 `**#**`(匹配零个或多个单词)通配符。这是最灵活的类型。 - **Headers Exchange**:不依赖 Routing Key,而是根据消息的 **headers** 属性进行匹配。 1. **投递消息**:将消息的副本发送到所有匹配的队列中。如果没有任何队列匹配,消息的行为取决于 Exchange 的配置(可以丢弃或返回给生产者)。  #### 消费者工作原理 消费者的核心任务是:**从指定的队列中获取消息,并可靠地处理它**。 1. **建立连接与通道**:与生产者相同,消费者也需要建立连接(Connection)和通道(Channel)。 2. **声明队列与绑定**: - - 为了健壮性,消费者通常会尝试声明它要消费的队列(如果队列不存在则创建)。 - 如果队列还未绑定到 Exchange,消费者也会执行绑定操作。 1. **订阅消息**: - - 消费者通过 `**channel.basicConsume()**` 方法向 Broker **订阅**一个队列。 - 这相当于告诉 Broker:“请持续将这个队列里的消息推送给我。” 这就是**推模式**,也是最常用的模式。RocketMQ 是长轮询拉取模式。 - TCP层面的心跳检测, AMQP协议内置心跳(`heartbeat` 参数),连接断开会立即检测到。 而不是应用层用ping来做心跳检测。 - brocker维护`queue → consumer` 映射表 ,能看到消费者的活性 1. **接收并处理消息**: - - 当队列中有新消息时,Broker 会通过通道**主动推送**给消费者。 - 消费者的回调函数被触发,接收到消息并开始执行业务逻辑。 1. **发送确认**: - - 消息处理完成后,消费者必须向 Broker 发送一个**确认**。 - **手动 ACK**:通过 `**channel.basicAck()**` 显式确认。这是**推荐做法**,可以确保消息在处理失败时不会被丢失(可以重新入队或进入死信队列)。 - **自动 ACK**:消息一发送给消费者就自动确认。性能好,但容易在消费者处理失败时丢失消息。  #### 死信队列工作原理 死信队列是 RabbitMQ **容错和监控**机制的核心组成部分。 死信队列本质上是一个普通的队列,从逻辑上它被用来存储“死掉”的消息。这里的“死”不是指消息本身损坏,而是指**消息因为某些原因无法被正常消费**。 一条消息在以下三种情况下会成为“死信”: 1. **消息被拒绝**:消费者使用 `**basicNack**` 或 `**basicReject**`,并且设置了 `**requeue=false**`(不重新入队)。 2. **消息过期**:消息在队列中存活时间超过了设置的 TTL(下面会讲)。 3. **队列已满**:队列达到了设置的最大长度,无法再存入新消息。 死信队列本身就是一个普通队列,它的特殊之处在于它是通过**死信交换机** 配置而来的。 **死信交换机(DLX, Dead Letter Exchange)** 是一种特殊的交换机,用来接收那些在正常队列中无法被消费的“异常消息”。 当消息在一个队列中变成了“死信”,RabbitMQ 会自动把这条消息转发到对应的 **死信交换机**,再由 DLX 路由到一个新的队列(称为 **死信队列 Dead Letter Queue, DLQ**) 在RabbitMQ中落地实现很简单,其实就是普通交换机什么也不用改,交换机名字用dlx.xx开头比较合理,绑定的队列也用dlx.xx开头,特殊之处在于正常队列配置一个map参数,指定死信交换机,指定死信路由键就行。 死信交换机: ```bash // channel.exchangeDeclare 就声明了一个交换机,名字叫做 dlx.exchange channel.exchangeDeclare("dlx.exchange", "direct"); // 声明一个队列,名字叫做dlx.queue,意为死信队列,本质上是普通队列 channel.queueDeclare("dlx.queue", true, false, false, null); // 队列绑定死信交换机 channel.queueBind("dlx.queue", "dlx.exchange", "dlx.key"); ``` 正常交换机指定死信交换机,这一步才是真正给死信交换机赋予含义: ```java Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", "dlx.exchange"); // 指定死信交换机 args.put("x-dead-letter-routing-key", "dlx.key"); // 指定死信路由键(可选) channel.exchangeDeclare("normal.exchange", "direct"); // 最后一个参数不是null就代表有死信队列 channel.queueDeclare("normal.queue", true, false, false, args); channel.queueBind("normal.queue", "normal.exchange", "normal.key"); ``` 工作原理图:  **死信队列使用场景**:降级处理、延时队列(别把死信队列看成故障的队列,就是个正常的队列) #### TTL(TIME-To-Live) **生存时间** TTL,即**生存时间**。它可以为消息或队列设置一个过期时间,超过这个时间后,消息就会“死亡”。T TL 是 RabbitMQ **消息生命周期管理**的基础功能。 **TTL 的两种设置方式** 1. **队列 TTL**: - - 在创建队列时,设置 `**x-message-ttl**` 参数(单位:毫秒)。 - **效果**:所有进入该队列的消息都会继承这个过期时间。 - **特点**:一旦设置,队列中所有消息的过期时间都一样,无法为单条消息定制。 1. **消息 TTL**: - - 在发送每条消息时,设置 `**expiration**` 属性(单位:毫秒)。 - **效果**:只有这条消息拥有独立的过期时间。 - **特点**:可以为每条消息设置不同的过期时间,更加灵活。 **TTL 的一个重要“坑”** **消息过期后,并不会立即从队列中删除!** RabbitMQ 只有在两种情况下才会处理过期消息: 1. 消息**即将被消费者消费**时(到达队列头部),没到达队列头部那就不会清理,哪怕过期了也是存在。 2. RabbitMQ 的一个**惰性扫描线程**定期检查队列。 这意味着,如果队列积压严重,一条已经过期的消息可能还会在队列中存活一段时间,直到它被扫描到或到达队首。如果配置了 DLX,过期消息在被处理时就会变成死信。 工作原理图:  #### 优先队列工作原理 一个支持消息优先级的队列。当队列中有消息积压时,高优先级的消息会**优先于**低优先级的消息被消费者获取。 优先级队列让 RabbitMQ 具备了**按重要性处理消息**的能力。 工作原理: 1. **启用优先级**:在创建队列时,必须设置 `**x-max-priority**` 参数(这就成为**优先队列**了),定义该队列支持的最大优先级(例如 5,表示优先级范围是 0-5)。 2. **设置消息优先级**:发送消息时,在消息属性中设置 `**priority**` 字段(值必须在队列支持的最大优先级范围内)。 3. **内部排序**:当消息进入队列时,RabbitMQ 并不是简单地追加到队尾。它会将高优先级的消息**插入到队列中较低优先级消息的前面**。 4. **消费顺序**:消费者总是从队列头部获取消息,因此高优先级的消息会先被消费。 优先级队列结构图:  #### 镜像队列工作原理 镜像队列:一个普通的队列,其内容可以被**实时复制**到一个或多个其他 Broker 节点上。形成“一主多从”的架构。 - **主节点**:负责处理所有对队列的读写操作(生产者发送、消费者消费)。 - **从节点**:作为热备,会从主节点同步所有消息和状态,但不对外提供服务。 镜像队列是 RabbitMQ **传统的高可用(HA)解决方案**,用于保证在主节点故障时,队列中的消息不丢失,服务不中断。 ⚠️**已过时**:在 RabbitMQ 3.8+ 版本后,官方推荐使用功能更强大、设计更现代的**仲裁队列**来替代镜像队列。 工作原理: 1. **配置镜像策略**:管理员需要在 RabbitMQ 管理界面或通过命令行设置一个**镜像策略**。这个策略定义了哪些队列需要被镜像,以及镜像到哪些节点上。 2. **主从同步**:当一个队列被策略匹配后,它就成为镜像队列。 - - 所有生产者发送的消息,都会先写入主节点。 - 主节点会将消息**同步**给所有从节点。 - 只有当**所有从节点**都确认收到后,主节点才会向生产者发送确认(ACK)。这保证了数据的强一致性。 1. **故障转移**: - - 如果主节点宕机,RabbitMQ 集群会自动从从节点中**选举一个新的主节点**。 - 原来的从节点升级为主节点,开始对外提供服务。 - 这个过程对生产者和消费者是**透明**的,它们会自动重连到新的主节点。 #### 仲裁队列工作原理 仲裁(Quorum)队列是 RabbitMQ **3.8 版本后推出的新一代高可用队列类型**。它基于 **Raft 共识算法**实现,旨在提供一个**数据安全、强一致、配置简单**的队列解决方案,是官方推荐的镜像队列替代品。 目标是实现**高可用**的**消息队列**架构。 “Quorum”一词意为“法定人数”,这暗示了它的核心机制:**需要大多数节点同意**,操作才能成功。 **工作原理:基于 Raft 共识算法** 仲裁队列的核心是 Raft 算法在消息队列场景下的实现。我们可以把它想象成一个**民主委员会**来管理队列。 一个仲裁队列通常部署在奇数个 Broker 节点上(最少 3 个),每个节点上的队列副本扮演一个角色: - **Leader(领导者)**:**唯一**的领导者,负责处理所有来自生产者和消费者的请求(读写操作)。 - **Follower(跟随者)**:普通的委员,不对外提供服务。它们的主要工作是**复制** Leader 的所有操作,并投票。 - **Learner(学习者)**:观察员角色,只同步日志,不参与投票。用于在不影响性能的情况下增加副本数,用于灾备或读取扩展。 - **消息写入流程(核心)** 这是理解仲裁队列的关键,它保证了数据的强一致性。 1. **所有写操作必须通过 Leader**。 2. **消息不是写入就成功**,而是需要**超过半数**的节点(包括 Leader 自己)都确认写入日志后,才算真正“提交”。 3. **只有已提交的消息才能被消费者消费**。这确保了即使 Leader 宕机,新选举出的 Leader 也一定拥有所有已提交的消息,**数据零丢失**。  - **故障转移流程** **心跳检测**:Leader 会定期向所有 Follower 发送心跳。 **触发选举**:如果 Follower 在一段时间内没有收到 Leader 的心跳,它就认为 Leader 宕机,于是发起选举,将自己转为 Candidate(候选人)。 **投票选举**:Candidate 向其他节点请求投票。获得**超过半数**选票的节点成为新的 Leader。 **数据恢复**:新 Leader 上拥有所有已提交的消息,可以立即对外提供服务,整个过程**自动完成**,对客户端透明。 #### 延迟队列工作原理 延迟队列(延迟消息):一个队列,其中的消息不会立即被消费者消费,而是在等待一段指定的时间后,才变成可消费状态。 延迟队列是一个非常常见的需求,但**RabbitMQ 本身不直接提供延迟队列功能**。我们需要通过插件或巧妙的组合来实现。 **实现方式一:TTL + DLQ(经典组合)** 这是最常用、最巧妙的实现方式,不依赖任何插件。 1. **架构设计**: - - 死信交换机:创建一个**业务交换机**和**业务队列**(不设置 TTL)。 - 延迟交换机:创建一个**延迟交换机**和**延迟队列**(设置 TTL)。 - 将延迟队列绑定到延迟交换机,并设置**死信交换机为业务交换机。** 1. **工作流程**: 1. 1. 生产者发送消息到**延迟交换机**,并设置 TTL。 2. 延迟交换机将消息路由到**延迟队列**。 3. 消息在延迟队列中等待,直到过期,期间不会有任何消费者拉取。 4. 消息过期后,成为死信,被发送到**死信交换机**(即业务交换机)。 5. 业务交换机根据路由规则,将消息路由到**业务队列**。 6. 消费者从业务队列中消费消息。 **实现方式二:延迟插件** RabbitMQ 官方提供了一个 `**rabbitmq_delayed_message_exchange**` 插件,提供了更原生、更优雅的实现。 1. **安装插件**:在所有 RabbitMQ 节点上安装并启用该插件。 2. **使用**:插件会提供一个新的交换机类型 `**x-delayed-message**`。 3. **工作流程**: 1. 1. 生产者发送消息到 `**x-delayed-message**` 类型的交换机。 2. 在消息头中添加 `**x-delay**` 属性,指定延迟时间(毫秒)。 3. 交换机收到消息后,**不会立即路由**,而是将其保存在一个内部的 Mnesia 表中。 4. 插件的后台定时器会定期检查这些消息。 5. 当消息的延迟时间到达后,交换机才会像普通交换机一样,将消息路由到目标队列。  #### 批量发送/批量确认工作原理 **批量发送**:**生产者**不是一次发送一条消息,而是将多条消息打包,一次性通过一个网络请求发送给 Broker。 **批量确认**:**消费者**不是处理完一条消息就发送一次 ACK,而是处理完一批消息后,发送一个 ACK,确认这批消息都已成功处理。 这是 RabbitMQ **性能优化**的关键手段,通过减少网络往返次数来提升吞吐量。 **批量发送原理** - **客户端实现**:批量发送主要是**客户端**的行为。RabbitMQ 的 Java 客户端(如 Spring AMQP)提供了 `**RabbitTemplate.convertAndSend**` 的批量版本。 - **网络效率**:将 N 条消息的 N 次网络请求,合并为 1 次网络请求,极大地减少了网络 RTT(往返时间)的开销。 - **权衡**:批量发送会增加客户端的内存使用,并且会带来一定的延迟(需要等一批消息凑齐或超时),空间换时间。 **批量确认原理** 1. **开启手动确认**:消费者必须关闭自动 ACK (`**autoAck=false**`)。 2. **处理消息**:消费者从队列中拉取一批消息(或 Broker 推送一批),并依次处理。 3. **发送批量 ACK**:当一批消息都处理成功后,消费者调用 `**channel.basicAck(deliveryTag, multiple=true)**`。 - - `**deliveryTag**`:这批消息中**最后一条**消息的标签。 - `**multiple=true**`:这是关键!它告诉 Broker:“请确认并删除**小于等于**这个 `**deliveryTag**` 的所有消息”。 1. **Broker 处理**:Broker 收到这个批量 ACK 后,会一次性将队列中这批已确认的消息全部清除。  ### RocketMQ 消息队列概念汇总 RocketMQ 非常适合应用在微服务架构中,经常作为微服务的消息中间件,所以下面的概念以微服务视角去思考、代入才会更好理解。 我们来看下面这张图,它囊括了 RocketMQ 相关的组件和角色。  核心概念一览表 | 概念 | 定义 | 对应 RabbitMQ 类比 | 说明 / 使用场景 | | -------------------------------- | ------------------------------------- | ------------------------------------------------- | --------------------------------------------- | | **生成者和消费者** | | | | | **Producer** | 消息生产者 | RabbitMQ Producer | 发送消息到 Topic / Queue | | **Producer Group** | 消息生产者群组 | | 事务,管理多个生产者 | | **Consumer** | 消息消费者 | RabbitMQ Consumer | 拉取或订阅消息 | | **Consumer Group** | 消费者组 | RabbitMQ Queue + 多消费者 | 一组消费者共享 Topic 消息,实现负载均衡或广播 | | **中间队列层概念** | | | | | **Broker** | 消息服务器 | RabbitMQ Broker | 消息存储和分发节点 | | **NameServer** | 元数据服务 | RabbitMQ 没有对应(Exchange/Binding+Cluster管理) | 集群路由和 Broker 元数据注册中心 | | **Topic** | 消息主题 | RabbitMQ Exchange(逻辑路由) | 消息的逻辑分类,Producer 发送消息到 Topic | | **Message Queue (队列)** | 物理队列 | RabbitMQ Queue | 存储消息的物理队列,由 Broker 管理 | | **commitLog** | RocketMQ 核心物理文件 | / | 存储所有消息的实际内容 | | **ConsumeQueue** | 队列的索引文件 | / | 一条记录映射一条commitLog 记录 | | **Message Queue Offset** | 队列偏移量 | RabbitMQ 没有明确概念 | 消费进度记录在 Broker 或 Consumer | | **消息层概念** | | | | | **Message** | 消息对象 | RabbitMQ Message | 包含 Body、Topic、Tag、Keys 等 | | **Tag** | 消息标签 | RabbitMQ RoutingKey | 消费端可做消息过滤 | | **Message Key** | 消息唯一标识 | RabbitMQ MessageId / header | 用于消息追踪和查找 | | **顺序消息** | 消息保证严格顺序消费 | RabbitMQ 默认 FIFO(Queue 内) | 消息队列分区内顺序处理 | | **RocketMQ 事务消息** | Producer 发送事务消息,最终提交或回滚 | RabbitMQ Transaction 或 Confirm 模式 | 保证分布式事务可靠性 | | **延迟消息 / Scheduled Message** | 消息延迟投递 | RabbitMQ TTL + 延迟插件 | 定时任务、延迟任务 | | **死信消息(DLQ)** | 消费失败或超过重试次数的消息 | RabbitMQ DLX | 异常消息处理 | 概念很多,最核心的是 commitLog 、Consume Queue 、 message key 、事务消息 #### 生产者和消费者 **1) Producer** : **消息的发送方**,负责把消息发送到 **Topic(逻辑分类),** 消息通过 Broker 进行存储,最终写入 CommitLog, Producer 可以选择不同的发送模式: - **同步发送(sync)**:等待 Broker ACK,保证消息可靠性 - **异步发送(async)**:回调方式,适合高吞吐量场景 - **单向发送(one-way)**:不等待 ACK,适合日志或监控消息 发送消息可以对消息指定: - **指定 Topic 和 Tag** ,用于消息分类和过滤 。 - **指定 Message Key** ,业务唯一标识,便于查询或幂等处理 - **事务消息支持** , 可以发送事务消息,配合 Broker 做事务半消息/回查机制 **2) Producer Group(生产者群组)** : 每一个生产者实例在创建时,一定要指定一个**生产者组**名,因此形成生产者组的概念。作用是 事务消息回查的关键标识 —— Producer 发送半事务消息,Brocker总是等不到提交指令,就会根据生成者组名回查生产者现在是什么状态。 **3) 消费者**: 消费者是一个独立运行的程序或进程(比如一个 Spring Boot 微服务实例)。它的任务是从 RocketMQ 的服务器(Broker)获取消息,并执行具体的业务逻辑(比如扣减库存、发送短信、生成报表等)。 **4) Consumer Group(消费者组):**在消费者创建时必须指定一个消费者组名字,消费者组也就因此形成。作用是**处理同一类业务逻辑**的消费者应用实例会被归为同一个组,增大这一类业务的消费吞吐量。 - 集群模式(点对点):消费组会**分摊消费**Brocker 中的 Topic 消息,一个消费者实例对接一个Topic下的队列。 - 广播模式: 消费者组内的每一个消费者,都会收到该主题下的**全量**消息 。 consumer - Brocker 传递数据的过程: 1. 消费者主动向 mq 发送一个长轮询请求 2. 如果有数据立即返回 3. 如果没有数据,挂起请求15秒(消费端指定**长轮询**的时间,挂起过程不占用线程),期间内有数据,唤醒请求返回数据。期间内没有数据,返回空响应(空转)。 ------ #### 中间层队列概念 **1)Broker:** 消息存储与转发器,负责接收生产者的消息,接收的消息持久化存储到commitLog,Broker提供给消费者消息拉取能力。生成者通过NameServer 找到 Brocker **2)NameServer** : 元数据注册中心,管理 **Broker 地址** 和 **Topic 路由**信息 , 可以由轻量级服务器作为技术实现, 每个 Broker 启动时注册到 NameServer ,汇报broker的topic信息到看板并维持心跳检测。NameServer 会注册生产者组和消费者组的信息, **Producer 和 Consumer** 找NameServer 拉取 **Topic 路由信息并缓存。** **3)Topic (主题)**: 一个逻辑概念,由若干队列和Broker 绑定 topic 名字实现。 Topic 内部会被拆分成若干个 **队列**, 每个队列独立存储消息,形成并行的存储和消费结构。 Producer 和 Consumer 只需约定 Topic, 不关心彼此 ,即可互相通信, 实现解耦**。**Producer 通过 NameServer 查询 Topic 的分区信息 。 **4)message queue** : 存储消息的队列容器,消息在队列内有顺序编号(MessageQueue Offset),用于顺序消费 ,负责对topic **5)commitLog** : 是 RocketMQ 的 **核心物理存储文件**, 存储的就是 **消息的物理内容,** 它按顺序 顺序追加(append) 所有消息,无论消息属于哪个 Topic 或 MessageQueue, 所以整个RocketMQ 仅有一个 commitLog . commitlog = 消息内容 + 消息属性 + 系统元数据 + 长度信息 **6)Consume Queue** : 队列的元信息 , ConsumeQueue 是索引文件 ,每条记录对应 MessageQueue 的一个 Offset,记录了该逻辑消息在 CommitLog 中的物理偏移,Consumer 通过它快速找到消息 ,内部结构 : ```bash ConsumeQueue[0] = <CommitLog offset, size=128, tag> ConsumeQueue[1] = <CommitLog offset, size=128, tag> ConsumeQueue[2] = <CommitLog offset, size=128, tag> 这里的[0][1][2] 便是 Message Queue Offset ``` **7)Message Queue Offset :** 每一个消息中都有一个顺序编号,称为Offset , 这是 **物理队列内消息的唯一序号**,编号从0单调递增,假如队列内3条消息消费完了,再来一条新消息,从4开始编号, 不会重用 ,永远递增。 - RocketMQ 崩溃后如何恢复编号? 通过 commitLog + ConsumeQueue - 重启后, Broker 会 **读取 ConsumeQueue 文件尾部**, 找到最后一条记录 `<CommitLog offset, 消息长度, Tag>`, 逻辑序号 = 该记录在文件中的顺序位置 , 即使消息已经被 Consumer 消费,ConsumeQueue 索引仍存在,所以可以恢复最大 Offset . - **集群模式下,由 Broker 集中管理****Message Queue Offset,**确保一致性和可靠性。 - **广播模式下,由消费者各自本地管理****Message Queue Offset,**简单高效。 Rocket MQ 的 commitLog 是关键设计, commitLog 会保存已经消费的消息,不删任何消息,这样便支持重复消费,支持了重复消费,那么 同一个 Topic 就可以被多个 Consumer Group 消费 ,并以 commitLog 的设计为基础设计了事务消息。 注意: RocketMQ 并不是无限保留 CommitLog 消息 ,**消费完成并过期的消息**,会由 **定期清理(Clean-up)机制** 删除 ------ #### 消息层概念 **1) Message(消息):** 是 RocketMQ 中 **最基本的数据单元,** 表示 业务系统需要传递的事件或数据, 主要包含三个部分: - 主题(Topic) : 消息逻辑分类,消费者按 Topic 消费 - 消息体(Body) : 业务数据,通常是字节数组(例如 JSON、字符串) - 消息属性(Properties) : 元数据,可选,用于 Tag、Key、延迟、事务、过滤等 详细字段: | **字段** | **说明** | **用途** | | -------------- | ------------ | -------------------------------- | | Topic | 消息主题 | 消费者订阅、分类 | | Body | 消息内容 | 实际业务数据 | | Tag | 消息标签 | 消费端可选择性订阅,用于消息过滤 | | messge Key | 消息唯一标识 | 幂等、查询、追踪 | | DelayTimeLevel | 延迟级别 | 延迟消息发送 | | BornTimestamp | 消息产生时间 | 消息顺序/延迟计算 | | Properties | 扩展属性 | 自定义字段,可用于业务逻辑 | **2) Tag (标签)**: 理解成消息的一个字段,消息的元信息,动态的随消息生成。多个topic可以有相同的tag,不影响。 **3) message Key**: 是消息的唯一标识符或 业务唯一标识**,** 它是由生产者业务系统设置的标识 , 通常用于业务系统追踪或查询消息 , 可用于 **消息查询、幂等、日志追踪、统计,** Message Key 尽量唯一,但也可以多个消息共用一个 Key(例如同一个订单号) RocketMQ 存储消息时,Key 仅作为 **消息属性**,实际消息定位仍靠 **CommitLog + ConsumeQueue Offset** Message Key 可以唯一,也可以不唯一 - **唯一 Key** - - 适合 **业务系统要求严格幂等或精确查询**的场景 - 例如:订单号、支付流水号 - 查询时通过 Key 能唯一定位消息 - **非唯一 Key** - - 同一个业务对象可能发送多条消息,但 Key 可以相同 - 例如:同一个用户多次操作,Key 是用户 ID - 查询时可以返回多条消息 **4) 顺序消息:** 消息的发送顺序与消费顺序保持一致。 RocketMQ实现顺序消息的关键在于 **“同一业务的消息始终进入同一个 MessageQueue”**,然后在消费端 **同一时刻只由一个线程顺序消费该队列的消息**。 整个过程依赖队列内天然的FIFO. ------ ##### 1. 事务消息 **事务消息**:它确保 消息的发送 与 本地事务的执行 要么都成功,要么都失败,从而保证分布式场景下的数据一致性 本地事务指的是生产者服务开启了一个事务,将业务操作和发生消息都纳入到事务步骤中,同成功同失败。  **半消息状态**:半消息会被标记一个特殊属性 `PROPERTY_TRANSACTION_PREPARED`,设置为 `"true"`,事务消息的独有特征,半消息不会直接发送到指定的业务topic中,RocketMQ 采用了 topic 隔离策略来存储半消息,这个 topic 也是个普通的topic, 但是消费者不能消费,便具有了半消息的语义。 **存半消息的 topic** :RocketMQ 内部名为 `RMQ_SYS_TRANS_HALF_TOPIC`的 Topic ,此 Topic 是内置的、天生的,作用是业务消费者无法订阅它。 **操作消息(OP消息)**:这是一个 RMQ_SYS_TRANS_OP_HALF_TOPIC,半消息提交或回滚的操作也叫做一条信息,存储在这里。因为 RocketMQ 的存储机制是基于 commitlog 顺序写 + ConsumeQueue 索引文件。 **消息一旦写入,就不会修改(Append-only 模式)**。 Broker 不能直接更新 Half Message 的状态, 所以采用“写一条对应操作记录”的方式来表达状态变化 ,**这和数据库的“WAL 日志(Write-Ahead Log)”很像**: 不改原数据,而是追加一条变更操作日志。 **事务回查**:rocketMQ 向 producer 询问当前的事务状态,producer 会检查本地事务是成功还是失败 **定时任务**: Broker 定期扫描 Half Queue,找出需要回查的消息 **Producer 回查接口** : 生产者实现业务逻辑检查点 在分布式系统中,最难的问题之一就是 **跨系统的一致性**,既 消息系统(MQ)和数据库事务无法保证**原子性(Atomicity)**。 RocketMQ 的 **事务消息** 正是为了解决这一点而生的。 它提供了一个“**两阶段提交 + 回查机制**”的模型,保证两者的最终一致性。 缺点:它只保证**最终一致性**,且中间状态(如支付成功但下游服务异常)在短暂时间内可能不一致,需要业务方自行处理。 Kafka、RabbitMQ 等常见 MQ 原生并没有类似机制 ,而这正是 RocketMQ 的亮点之一,国产之光! ------ ##### 2. 延迟消息 **延迟消息(Scheduled Message)** 是指消息在发送到 Broker 后,不会立即被投递给消费者,而是**延迟一段时间后**再被消费。 实现原理: - 生产者形成延迟消息时 , 带上 delayLevel,Broker 收到后 **不会立即写入真实 Topic** - Broker 写入“延迟队列”(定时队列) , 延迟队列的 Topic 固定为 `SCHEDULE_TOPIC_XXXX` - Broker 定时任务扫描到期消息 , 到期后将消息重新写入原始 Topic(即用户真正的 Topic) - 消费者收到消息 , 从原 Topic 消费到这条延迟后转发的消息 这跟RabbitMQ 的 **TTL + DLQ 实现延迟队列思想一致。** ------ ##### 3. 死信消息 **死信消息**:指在消息消费失败,且达到**最大重试次数**后(默认16次)仍未成功,RocketMQ 会将这类消息转移到特殊的队列中进行隔离,这些消息被称为死信消息,存储死信消息的队列就是死信队列 **死信队列**:它是一个**特殊的 Topic**,名称通常为 `**%DLQ% + 消费组ID**`。每个消费组都有其对应的死信队列,用于存储该消费组内所有 Topic 的死信消息 死信队列是 RocketMQ **消息可靠性设计的最后一道防线**,旨在**隔离问题消息**,防止其影响正常消息的处理流程,并提供一个集中的地方供监控和人工干预。 使用场景:死信队列主要用于**处理消费失败且无法通过重试恢复的消息** - **订单处理失败**:订单创建、支付、库存扣减等操作失败的消息,可以进入死信队列,后续通过人工干预或自动程序进行补偿 1. **异常消息处理**:由于数据格式错误、业务逻辑异常等原因无法消费的消息,进入死信队列,供后续排查和修复。 2. **超时消息处理**:设置了 TTL 但未能在规定时间内消费的消息。 3. **监控与告警**:监控死信队列可以及时发现系统异常,是系统健康度的重要指标。 ------ #### 生产者工作原理 单生产者工作原理  1. **创建并启动 Producer** - - 程序创建 `**DefaultMQProducer**` 实例。 - 设置 NameServer 地址、Producer Group 名称。 - 调用 `**start()**` 方法启动 Producer。 1. **获取路由信息** - - Producer 启动后,会从**一个 NameServer**(随机选择)拉取 Topic 的路由信息。 - 路由信息包含:该 Topic 分布在哪些 Broker 上,每个 Broker 上有哪些 MessageQueue(队列)。 - Producer 会**定时(默认30秒)** 更新本地路由缓存。 1. **选择 MessageQueue** - - Producer 根据负载均衡策略(如轮询 `**RoundRobin**`),从获取到的 MessageQueue 列表中选择一个队列。 - 如果选择了写失败的队列,会进行重试,并可能暂时规避该队列。 1. **构建并发送消息** - - 将消息体、Topic、Tag、Key 等信息封装成消息对象。 - 对消息进行序列化、压缩等处理。 - 通过 Netty 客户端向选定的 Broker 发送消息。 1. **Broker 处理消息** - - Broker 接收消息后,将其**顺序写入 CommitLog 文件**。 - 写入成功后,返回一个包含消息物理偏移量(CommitLog Offset)的 ACK 给 Producer。 1. **Broker 异步构建索引** - - Broker 的后台线程会**异步**地将 CommitLog 中的消息位置、大小等信息分发到对应的 ConsumeQueue 文件。 - 如果消息带有 Key,还会构建 IndexFile 索引,以支持按 Key 查询。 1. **Producer 接收 ACK** - - Producer 收到 Broker 返回的 `**SendResult**`,状态为 `**SEND_OK**`,表示消息发送成功。  重点说明: - **异步构建索引**:消息写入 CommitLog 和构建 ConsumeQueue 索引是**异步**的。这极大地提升了 Broker 的吞吐性能,因为写入 CommitLog 是顺序写,非常快,而构建索引可以稍后批量处理。 - **高可用设计**:如果发送失败,Producer 会自动重试(默认2次),并可能选择其他 Broker 上的队列,实现了容错。 **生产者群组工作原理**  **1. 启动与路由发现** - **拉取路由**:当 Producer Group 中的任何一个实例启动时,它会**定期向 NameServer 拉取**它所关心的 Topic 的路由信息(包括该 Topic 分布在哪些 Broker 上,以及每个 Broker 上的 MessageQueue 列表)。 - **心跳机制**:Producer 实例会与 NameServer 和 Broker 保持心跳,但这主要是为了让它们感知自己的存活状态,而非为了“注册”。 - **无状态设计**:因此,Producer 实例本身是**无状态的**,Broker 也只关心消息内容,不关心消息具体来自哪个实例。这使得 Producer Group 内的实例可以任意增减,实现弹性伸缩。 **2. 消息发送(正常运行)** 在发送消息时,Producer Group 的工作机制体现了其高扩展性: - **独立发送**:每个 Producer 实例都**独立工作**,根据其本地缓存的路由信息,通过负载均衡策略(如轮询)从 Topic 的多个 MessageQueue 中选择一个。 - **并行处理**:然后,它将消息直接发送给该 MessageQueue 所在的 Broker。多个实例可以并行地向不同的 Broker 或不同的 Queue 发送消息,共同承担发送压力,实现了负载均衡。 **3. 事务消息回查(容错机制)** 在事务消息场景下,Producer Group 的设计提供了关键的容错能力,这也是其核心价值之一: - **问题场景**:如果某个 Producer 实例在发送“半消息”后、提交最终状态(COMMIT/ROLLBACK)前突然宕机,Broker 上就会留下一个状态不明的“孤儿”半消息。 - **回查机制**:此时,Broker **不会主动“通知”** 同组的其他实例。正确的流程是,Broker 会**主动发起回查**。它会根据半消息中记录的 **Producer Group 名称**,向 NameServer 查询该 Group 下**任意一个存活的 Producer 实例地址**。 - **状态恢复**:然后,Broker 向这个被选中的实例发起回查请求,由它代表整个 Group 去查询本地事务状态(如查询数据库),并将结果(COMMIT 或 ROLLBACK)返回给 Broker。Broker 根据这个结果完成最终的事务提交或回滚。  ------ #### 消费者工作原理 在深入两种模式之前,先理解几个贯穿始终的核心概念: - **Consumer Group(消费者组)**:一个逻辑概念,由多个消费者实例组成。它们共同消费一个或多个 Topic,是实现负载均衡和容错的基本单元。理解消费者实例和队列对应关系很重要。 - **Rebalance(重平衡)**:一个动态过程,当消费者组内的实例发生变化(加入、离开)或 Topic 的队列数量变化时,Broker 会重新分配队列与消费者的对应关系。 - **Offset 管理**:记录消费进度。这是区分两种模式的关键所在。 **1)集群模式** 集群模式是**默认且最常用**的模式,旨在通过横向扩展来提升消费能力。 1. **队列分配**:一个 Topic 的多个 MessageQueue 会被**分配**给消费者组内的不同实例。**一个队列在同一时间只会被组内的一个消费者实例消费**。 2. **负载均衡**:通过 Rebalance 机制,如果消费者实例数量多于队列数量,多余的实例会闲置;如果实例数量少于队列数量,一个实例会消费多个队列。这实现了消费任务的负载均衡。 3. **故障转移**:如果一个消费者实例宕机,它负责的队列会被 Rebalance 机制重新分配给其他存活的实例,确保消息消费不中断。 4. **Offset 存储**:由于消费进度需要被组内所有实例共享(尤其是在故障转移时),消费进度(Offset)**集中存储在 Broker 端**。  **2)广播模式** 广播模式适用于需要将同一条消息推送给所有下游消费者的场景,如配置更新、状态通知等。 1. **全局消费**:消费者组内的**每个实例都会消费 Topic 下所有 MessageQueue 的所有消息**。消息不会被分摊,而是被广播。 2. **无 Rebalance 分配**:由于每个消费者都要消费所有队列,因此不存在队列分配的 Rebalance 过程。每个消费者独立地订阅所有队列。 3. **独立消费**:各消费者实例的消费进度互不影响,一个实例消费慢或失败,不影响其他实例。 4. **Offset 存储**:由于每个消费者的消费进度都是独立的,不需要与其他实例同步,因此消费进度(Offset)**存储在消费者本地**。  两种模式对比: | 特性维度 | 集群模式 | 广播模式 | | --------------- | ---------------------------------------------------- | -------------------------------------------------------- | | **消费关系** | 一个队列**只被**一个消费者实例消费 | 一个队列**被所有**消费者实例消费 | | **负载均衡** | **支持**,通过 Rebalance 在组内分摊消费任务 | **不支持**,每个实例承担全部消费负载 | | **故障转移** | **支持**,宕机实例的队列会被其他实例接管 | **不支持**,实例间完全独立,互不影响 | | **Offset 存储** | **Broker 端**集中存储 | **消费者本地**独立存储 | | **适用场景** | 高吞吐量、分摊消费压力的业务(如订单处理、日志分析) | 配置下发、状态通知、所有客户端需要同步收到相同消息的场景 | ------ ### kafka 消息队列概念汇总 | 类别 | 概念 | 作用 | 类比与说明 | | -------------------- | -------------------------------- | ------------------------------------------------------------ | ------------------------------------------------------------ | | **核心组件** | **Producer** | 消息生产者,负责发送消息到 Kafka Broker。 | 消息的来源,可以是网站前端、后端服务、日志采集器等。 | | | **Consumer** | 消费者,从 Kafka 读取消息。 | 消息的终点,负责处理消息,如数据分析、业务处理等。 | | | **Broker** | Kafka 服务器实例,负责存储和转发消息。 | 一个 Kafka 集群由多个 Broker 组成,每个 Broker 都是一个独立的节点。 | | | **Topic** | 消息分类主题,是逻辑上的消息分类。 | 类似数据库中的表名,用于区分不同类型的消息(如订单、用户行为)。 | | | **Partition(queue)** | Topic 下的分区,是真正存储消息的地方。partition 就是**队列** | 一个 Topic 可以分为多个 Partition,分布在不同 Broker 上,实现水平扩展和并行处理。**类比 RocketMQ 的 MessageQueue**。 | | | **Offset** | 每条消息在 Partition 中的唯一序号。 | **Partition 内部**消息的指针,从 0 开始单调递增。**类比 RocketMQ 的 Message Queue Offset**。 | | | **Consumer Group** | 一组消费者,负责负载均衡消费。 | 一个 Topic 的消息可以被多个不同的 Consumer Group 订阅,每个 Group 都有独立的消费进度。 | | | **Leader / Follower** | Partition 的主从副本,保证高可用。 | 每个 Partition 有一个 Leader 负责读写,多个 Follower 负责同步数据。**类比 RocketMQ 的 Master/Slave**。 | | | **Zookeeper / Kafka Controller** | 管理 Broker、Topic、Partition 元数据。 | **旧版依赖 Zookeeper**。**新版(KRaft)用 Controller 代替**,Controller 是从 Broker 中选举出来的,简化了架构。**类比 RocketMQ 的 NameServer**。 | | **生产者端** | **acks** | 生产者可靠性核心配置,决定何时认为消息发送成功。 | `**acks=0**` (发完即成功), `**acks=1**` (Leader收到即成功), `**acks=all**` (所有副本收到才成功)。 | | | **幂等性生产者** | 保证单分区单会话内 Exactly-Once,防止重试导致消息重复。 | 开启后,Producer 为每条消息分配唯一ID,Broker 去重,确保消息不重复。 | | | **事务性生产者** | 保证跨分区、跨主题的原子写入,实现“要么都成功,要么都失败”。 | 结合幂等性和事务协调器,实现 Kafka 的 Exactly-Once 语义 (EOS)。 | | **存储与 Broker** | **Log Segment** | Partition 内的日志文件片段。 | 一个 Partition 的日志由多个 Log Segment 组成,方便日志滚动和清理。**类比 RocketMQ CommitLog 的切分文件**。 | | | **Log 结构** | Kafka 高性能的基石,Partition 是一个只追加的、有序的日志文件。 | 顺序写磁盘,避免了随机写的巨大开销,性能极高。 | | | **零拷贝** | 消费者高性能的关键,数据在内核态直接从磁盘发送到网卡。 | 避免了数据在用户空间和内核空间之间的多次拷贝,极大提升了吞吐量。 | | | **索引** | 加速消息查找,包括 Offset 索引和时间戳索引。 | 允许消费者根据 Offset 或时间戳快速定位消息,而不必遍历整个日志。 | | | **Retention Policy** | 消息保留策略,按时间或大小删除旧消息。 | Kafka 的清理机制,防止磁盘被占满,如保留7天或达到10GB大小。 | | **消费者端** | **Rebalance** | 消费者组内分区分配的动态过程,是负载均衡和故障转移的基础。 | 当组内消费者数量或 Topic 分区数量变化时触发,期间消费会短暂停止。 | | | **提交机制** | 决定消费语义的关键,消费者如何告知 Broker 自己的消费进度。 | 分为自动提交、手动同步提交(可靠但慢)、手动异步提交(快但需处理回调)。 | | | **消费语义** | 消息传递的保证级别。 | At-Most-Once (最多一次), At-Least-Once (至少一次,默认), Exactly-Once (精确一次)。 | | | **位移主题** | 消费位移的存储位置。 | Kafka 将消费组的 Offset 信息存储在内部的 Topic `**__consumer_offsets**` 中,而非 Broker 文件。 | | **集群协调与高可用** | **ISR (In-Sync Replica)** | 副本同步集合,确保消息可靠复制。 | 与 Leader 保持同步的 Follower 集合。acks=all 时,Leader 必须等待 ISR 中所有副本确认。**类比 RocketMQ 的 Master-Slave 同步机制**。 | | | **KRaft 模式** | Kafka 的未来方向,移除 Zookeeper 依赖。 | 使用 Raft 协议在 Broker 内部选举 Controller,简化部署和运维,提升性能。 | | | **Cooperative Rebalance** | 更平滑的 Rebalance 协议,解决 Rebalance 时的“Stop-the-World”问题。 | 分区可以增量地、分批地重新分配,消费中断更短,体验更平滑。 | | **高级特性** | **消息压缩** | 节省网络带宽和磁盘空间。 | Producer 可在发送前对消息进行压缩(如 Gzip, Snappy),Broker 和 Consumer 自动解压。 | | | **安全机制** | 保障数据安全。 | 包括 SSL/TLS (加密传输)、SASL (身份认证)、ACL (权限控制)。 | **kafka VS RocketMQ 能力对比** | 特性 | Kafka | RocketMQ | | ---------------- | ---------------------------------- | ---------------------------------- | | 系统类型 | 分布式消息中间件 | 分布式消息中间件 | | 核心功能 | 异步解耦、削峰填谷、异步通信 | 异步解耦、削峰填谷、事务消息 | | 典型角色 | Producer / Broker / Consumer | Producer / Broker / Consumer | | 数据模型 | Topic → Partition → Offset | Topic → MessageQueue → Offset | | 消费模式 | 基于 **Pull(拉取)** | 基于 **Pull(主动拉)** | | 消息投递 | “至少一次” (At least once) | “至少一次”,支持事务一致性 | | 支持事务 | 有(Kafka Transaction) | 有(RocketMQ Transaction Message) | | 是否支持本地事务 | ❌ 否。Kafka 事务只管消息投递原子性 | ✅ 是。业务代码可以和消息发送绑定 | | 存储方式 | Partition 文件(日志结构) | CommitLog + ConsumeQueue | | 持久化机制 | Append-only + LogSegment | CommitLog 顺序写 | | Offset 管理 | Consumer 自管(可提交) | Broker 统一维护(集群模式) | | 路由发现 | Controller(或 ZooKeeper) | NameServer | | 高可用 | ISR 同步副本机制 | 主从 + 同步复制机制 | | 消息清理 | TTL 或磁盘大小策略 | 定期清理 CommitLog | | 应用方向 | 大数据、日志流处理、实时分析 | 分布式系统、事务一致性、业务异步化 | 本地事务: 不是消息队列的事务,而是 **你自己服务内部** 的那部分业务逻辑(mysql,Spring事务等)。 RocketMQ 可以做到 让 **本地事务** 和 **消息发送** 的结果保持一致。 Kafka做不到。 ## 三个消息队列对比 | 特性维度 | RabbitMQ | RocketMQ | Kafka | | ---------------- | ---------------------------------------------- | ------------------------------------------------ | ---------------------------------------------- | | **架构模型** | **中心化 Broker** (AMQP 协议) | **分布式消息平台** | **分布式日志系统** | | **核心设计目标** | **智能消息路由**和**可靠传递** | **金融级可靠性**和**海量消息**的强一致性 | **高吞吐量**的**实时数据流处理** | | **消息模型** | **Exchange -> Queue** 模型,路由灵活 | **Topic -> MessageQueue** 模型,借鉴 Kafka | **Topic -> Partition** 模型,简单直接 | | **性能吞吐** | 中等(万级/秒) | **高**(十万级/秒) | **极高**(百万级/秒) | | **可靠性** | 高(通过 ACK、持久化、镜像队列) | **极高**(同步/异步刷盘、事务消息、主从复制) | 高(通过多副本 ISR 机制) | | **顺序性** | **队列内有序**(多消费者会打乱) | **队列内有序**(严格顺序消息) | **分区内有序**(全局无序) | | **消息回溯** | 不支持(需插件) | **原生支持**(按 Offset 或时间) | **原生支持**(按 Offset 或时间) | | **事务消息** | **不支持**(可通过 AMQP 事务模拟,性能差) | **原生支持**(两阶段提交+状态回查) | **不支持**(需借助外部系统,如 Kafka Streams) | | **生态系统** | 成熟,支持多种协议 | 主要在 Java 生态,阿里系生态完善 | **极其庞大**(流处理、大数据生态) | | **运维复杂度** | **较低**(集群搭建相对简单) | **中等**(Java 技术栈,对国内团队友好) | **较高**(依赖 Zookeeper,配置复杂) | | **典型应用场景** | **微服务通信**、**业务事件驱动**、**任务队列** | **订单交易**、**金融支付**、**大规模互联网应用** | **日志收集**、**大数据管道**、**实时监控** | 哪个消息队列应用最广泛? 从全球范围来看,**Kafka 和 RabbitMQ 的应用广度远超 RocketMQ**。RocketMQ在国内应用广泛,毕竟起源于国内,专为 Java 生态而生。 1. **Kafka:大数据领域的绝对霸主** - - 在**大数据、日志聚合、流式计算**领域,Kafka 是事实上的标准。几乎所有的大数据组件(Flink, Spark, Storm)都与 Kafka 深度集成。 - 如果你的场景是处理海量的数据流,Kafka 几乎是首选。**在这个领域,它的应用最广泛。** 1. **RabbitMQ:企业级应用和微服务的常青树** - - 在**传统企业应用、微服务架构、业务系统集成**方面,RabbitMQ 凭借其成熟稳定、灵活的路由和丰富的协议支持,拥有巨大的用户基数。 - 很多公司的第一个消息队列就是 RabbitMQ,因为它**上手快**,能解决大部分业务解耦和异步通信问题。在通用业务集成领域,它的应用非常广泛。 1. **RocketMQ:亚洲和大型互联网公司的利器** - - RocketMQ 起源于阿里,在国内的**大型互联网公司、金融、电商**领域应用极为广泛。它经历了“双十一”等超大规模场景的考验。 - 它的优势在于对**金融级事务消息、顺序消息、海量消息堆积**等复杂场景的完美支持。 - 不过,在欧美市场,RocketMQ 的知名度和采用率远不及 Kafka 和 RabbitMQ。
Canal + RocketMQ监听MySQL环境搭建
  我们在项目中会引入缓存来提升系统性能。当一个请求进来,如果能够命中缓存,那么就直接返回缓存中的结果,不会去查询数据库。但引入缓存的同时,也就引入了与数据库之间不一致问题。 比如一个数据“AAA”存入了缓存,设置30分钟之后失效。那么这30分钟之内,查询到这条数据的请求拿到的就是“AAA”。但如果在第15分钟时,用户更新了这条数据,将它变成了“BBB”,这就导致了有15分钟的时间缓存与数据库之间不一致。   有一个方案能够解决这个问题:Canal监听MySQL的binlog,有数据变动就通过MQ发送数据更新的消息。后端接收到数据更新的消息,就可以及时的更新或者删除缓存,这样就避免了长时间缓存与数据库不一致。 ### 版本说明   MySQL 8.0.35   Canal 1.1.8   RocketMQ 4.9.7 ### 搭建步骤   这里牵扯到三方面内容,MySQL、Canal和RocketMQ。如果你是第一次搭建这个环境,那么我建议按照下面的步骤进行: #### 1 调整MySQL #### 2 安装并测试Canal的tcp模式   上述两步参考这篇文章:<a href="https://www.codefather.cn/post/1887431963256348674#heading-0">聚合搜索平台-SpringBoot整合Canal</a>。 #### 3 安装并测试RocketMQ   这一步网上的教程比较多,这里不再赘述。   这里要注意这样一个问题:RocketMQ是基于JDK8写的,并且Canal给MQ发消息时也需要调用MQ的jar包。所以如果你用的是高版本的JDK,那么可以修改Canal/RocketMQ启动文件中有关java的命令。这里以Canal启动文件为例。 ~~~ Windows: "D:\xxxxx\jdk1.8\bin\java.exe" %JAVA_OPTS% -classpath "%CLASSPATH%"…… ~~~ ~~~ Linux: ## set java path if [ -z "$JAVA" ] ; then JAVA=/.../jdk8/java/bin/java fi ALIBABA_JAVA=/.../jdk8/java/bin/java TAOBAO_JAVA=/.../jdk8/java/bin/java ~~~ ### Canal与RocketMQ进行整合 #### 修改Canal配置文件 canal.properties文件 ~~~ # 注意这个有的版本是RocketMQ,有的版本是rocketMQ,仔细看一下注释 canal.serverMode = rocketMQ rocketmq.producer.group = 你的groupName rocketmq.namesrv.addr = ip:端口 canal.mq.flatMessage = false canal.instance.filter.transaction.entry = true ~~~ example\instance.properties文件 ~~~ canal.mq.topic=你的topic # 这里稍后解释 canal.instance.filter.regex=数据库名\.表名1.*,数据库名\.表名.* canal.mq.dynamicTopic=true ~~~   关于canal.instance.filter.regex参数,有个大坑。我在搭这个环境的时候,tcp模式是没有问题的,但是切换到MQ模式之后,Canal就不发消息了,也不记录任何日志。这个bug我改了三天。   我在今天,先是把canal.mq.flatMessage改成了false。然后发现Canal发消息给MQ了,但是只发送事务开启和结束的消息EntryType.TRANSACTIONBEGIN和EntryType.TRANSACTIONEND。没有任何与数据变动有关的消息。   后来我也是在网上查到一篇文章,写的是canal.instance.filter.regex参数设置的不对的话,是会把所有表的数据变动全部过滤掉的(没错上面的bug基本上就是这个原因)。   那么tcp模式为什么是好的呢?如果你仔细观察Canal官方的示例代码,你会发现: ~~~ connector.subscribe(".*\\..*"); ~~~   我对于这个的理解就是,配置文件里的正则表达式,读取到程序里,可能就不是那个样子了。但是代码里面写的就不会有问题。tcp模式里有上面这行代码,就把问题给掩盖了。   这个正则表达式(.*\\..*)的意思就是监听所有库的所有表。但是在instance.properties这样写就不对,反而会忽略掉所有表,只发事务开启和关闭的消息。我也试过(.*\..*)和(.*),都不行。一个“\”是因为网上说多一个\是为了转义,但试过还是不行。   我建议就是(数据库名\.表名1.*,数据库名\.表名.*),这样写就正常 #### 后端代码测试 ~~~ import com.alibaba.otter.canal.client.CanalMessageDeserializer; import com.alibaba.otter.canal.protocol.CanalEntry; import com.alibaba.otter.canal.protocol.Message; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; import java.util.List; @Component @RocketMQMessageListener(topic = "topic名", consumerGroup = "group名") public class MessageConsumer implements RocketMQListener<MessageExt> { @Override public void onMessage(MessageExt message){ Message msg = CanalMessageDeserializer.deserializer(message.getBody()); List<CanalEntry.Entry> entries = msg.getEntries(); printEntry(entries); } private static void printEntry(List<CanalEntry.Entry> entrys) { for (CanalEntry.Entry entry : entrys) { if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN || entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) { continue; } CanalEntry.RowChange rowChage = null; try { rowChage = CanalEntry.RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException("ERROR ## parser of eromanga-event has an error , data:" + entry.toString(), e); } CanalEntry.EventType eventType = rowChage.getEventType(); System.out.println(String.format("================> binlog[%s:%s] , name[%s,%s] , eventType : %s", entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(), entry.getHeader().getSchemaName(), entry.getHeader().getTableName(), eventType)); for (CanalEntry.RowData rowData : rowChage.getRowDatasList()) { if (eventType == CanalEntry.EventType.DELETE) { printColumn(rowData.getBeforeColumnsList()); } else if (eventType == CanalEntry.EventType.INSERT) { printColumn(rowData.getAfterColumnsList()); } else { System.out.println("-------> before"); List<CanalEntry.Column> beforeColumnsList = rowData.getBeforeColumnsList(); printColumn(beforeColumnsList); System.out.println("-------> after"); printColumn(rowData.getAfterColumnsList()); } } } } private static void printColumn(List<CanalEntry.Column> columns) { for (CanalEntry.Column column : columns) { System.out.println(column.getName() + " : " + column.getValue() + " update=" + column.getUpdated()); } } } ~~~ 当你看到类似下面这样的输出,就表示成功了 ~~~ ================> binlog[xxx-bin.xxx:xxx] , name[数据库名,表名] , eventType : UPDATE (或者 INSERT 或者 DELETE) -------> before id : xxx update=false createTime : 2025-09-02 17:40:13 update=false editTime : 2025-09-02 17:40:13 update=false updateTime : 2025-09-13 01:12:21 update=false isDelete : 0 update=false -------> after id : xxx update=false createTime : 2025-09-02 17:40:13 update=false editTime : 2025-09-02 17:40:13 update=false updateTime : 2025-09-13 02:39:58 update=true isDelete : 0 update=false ~~~
深入浅出 RocketMQ 消息队列 笔记(13)
如何处理消息堆积问题 消息堆积的本质:消费速度<生产速度 常见消息堆积原因 1.瞬时流量 2.消费性能不足 3.机器不够 4.Bug 5.其他功能影响 瞬时流量 比如平常一分钟一条消息,瞬时流量来了,可能一秒钟就产生了上万条消息。 这里要对消息的瞬时性做分析,如果只是瞬时的,持续几分钟就好了,那就不需要做任何改动,毕竟消息队列就是负责削峰填谷的。如果持续时间长且消息量级过大,那堆积造成的影响还是需要重视。 这里举个例子,是yes哥之前遇到的一个情况,有七千万的历史数据需要清洗(消费),其中生产者花了12天时间将数据发给消息队列,但消费者一天只能消费200w条消息,所以乐观估计,消费者消费完所有历史数据需要至少35天。 这里会带来不少问题,如果消息队列用的是第三方提供的(比如阿里云),消息存储有个默认时间3天,如果消息存储3天还没被消费,那这个消息就会被删掉了(删掉的目的是为了控制存储成本),堆积未消费的消息丢失后,还需要生产者重新发送一遍,会很麻烦 还有一个问题,如果历史数据跟正常业务数据的处理流程是相同的,也就是他们的topic相同,那在broker产生堆积后,光顾着处理历史数据消息了,正常业务数据的处理也会因此而堵住了。  这里yes哥的解决方式是降低生产消息的频率,不持续发消息,让消息发送的频率略低于消费速率,给正常业务数据有消费的空间 消费性能问题 1.可以把循环里的查询拎出来放在外面批量查询,然后Map映射 2.单条插入/更新改成批量插入/更新 机器不够 消费性能上能优化还是优化了,但消息还是堆积,那就得水平扩展机器了,不过要注意,消费者与队列数要一直保持消费<队列的数量,否则会造成消费者空忙,不能很好的利用重平衡机制。  Bug 也有可能是消费逻辑有bug,导致消息堆积。消费失败需要重试,默认重试16次才会进入死信队列,相当于一条消息用到的流量被放大了16倍,资源都用来重试了,重试后还是失败。再如果消息是顺序消息,那后面的消息都处理不下去了。这种情况只能对消息队列做监控,发现报错后紧急发布版本,修复bug。进入死信队列的消息可以手动重新消费。 其他功能影响 其他业务功能抢夺了消息队列在使用的资源,比如消息消费的逻辑是扣减库存,其他业务功能也涉及到扣减库存,那么就有可能会争夺同一个商品的处理,也就是同一个数据库的行锁,导致消费速率降低。还有可能是其他功能的长事务导致消费等待,等等。如果平常正常,流量也没有变很大,有天突然堆积了,可以在这方面考虑问题原因。
深入浅出 RocketMQ 消息队列 笔记(12)
如何保证消息不重复 消息无法保证不重复,但可以保证消息被幂等消费,等同于仅消费一次 消息不重复解决方案 1.项目初期/新功能开始设计的时候,就考虑好幂等设计,满足消息的幂等消费 2.调整业务执行顺序,将影响大业务幂等的逻辑提前,提前终止重复消费 3.添加唯一性索引,如果有业务上的唯一索引(比如订单号)就不用加,反之就加一个唯一索引,重复消费插入时会报错的。如果不方便添加唯一索引,可以再开一张表存唯一索引,然后将这两张表放在同一个事务里 4.引入第三方中间件,比如redis,处理的时候用SETNX判断,已插入就直接返回,否则执行正常业务逻辑。不过这里有可能在redis刚存值的时候系统宕机了,导致消息被假处理,这个需要注意。
深入浅出 RocketMQ 消息队列 笔记(11)
保证消息不丢失的必要条件 生产者发送消息、生产者存储消息、消费者拉取消息,需要保证三大流程消息不丢失,缺一不可 生产者保证消息完整发送并存储至broker broker保证存储的消息不丢失 消费者保证拉取的消息一定被消费,即使重启了,也能确保未消费的消息继续消费 生产+发送消息流程 以下单为例,下单后增加积分,增加积分这个动作放在消息里实现,且要保证该消息一定发送成功。 这里可以引用TCP协议的ack请求确认机制。如果broker收到生产者推送来的消息,就返回ack给生产者,这样就能保证生产者发送消息这个阶段,消息不会丢失。如果生产者一直没收到ack,可能是网络或者其他原因,这时候就会进行重试,重试次数达到上限后抛异常。但这里抛异常不能影响主流程,所以需要对异常进行特殊处理。如果是同步发送消息,那就try-catch,如果是异步,那就在异步方法对应的onException做异常处理,这里可能要做一些补偿机制。 存储流程 broker在返回ack给生产者之前要确保消息已经成功存储了,RocketMQ的消息默认是异步刷盘,先刷到cache上,再等操作系统或定时刷盘任务把消息刷到磁盘上。如果此时断电了,消息就丢失了。所以为了保证消息不被丢失,这里可以把刷盘方式改成同步刷盘(flushDiskType=SYNC_FLUSH)。如果broker是集群,那也得保证主从broker的复制方式是同步复制,这样的话消息就更安全了。 消费流程 消费者消费后需要上报点位给broker,告诉broker已经消费到第几条消息了。如果消费者在处理消息的时候异步处理,然后直接告知broker消费成功,就很有可能在消费过程中报错,导致再次拉取消息的时候,从之前上报过的点位继续拉消息,消息就“丢失”了。所以这里要确保消费者真正消费完成消息后,再提交点位。
深入浅出 RocketMQ 消息队列 笔记(10)
RocketMQ的消息存储 用一个CommitLog所有分发给该broker的消息,多Topic混合存储 commitLog超过1G,会新起一个commitLog 每条消息存储到commitLog都会在consumeQueue生成一条记录,可以视为一个索引(稠密索引) Kafka的消息存储 Kafka在Topic下也分了多个队列来提高消费的并发度,但在Kafka不叫队列,叫分区partition Kafka的消息存储和RocketMQ略有不同,Kafka的消息存储以Partition为单位进行存储 每个Topic的每个分区都有自己的消息文件、索引文件和时间索引文件,他们的文件名相同,后缀名不同 文件名的命名规则是第一条消息 的offset,文件写满会新起一个文件 索引文件设计的和RocketMQ不同,Kafka是每隔几条消息再创建一条索引,节省了存储空间,能保存更多的索引,这样的索引叫稀疏索引 稀疏索引如何找到对应消息 通过offset找到对应的索引文件,通过二分遍历找到离消息最近的索引,再通过这个索引找到消息文件里此条消息的位置,再遍历消息文件找到目标消息 Kafka时间复杂度:O(log2n)+O(m),n为索引个数,m为稀疏程度 RocketMQ时间复杂度:O(1) 所以这里就需要权衡利弊了,Kafka是时间换空间,RocketMQ是空间换时间 虽然Kafka是按分区存储消息,但这样可能会引起 Kafka的存储设计对数据复制和迁移很友好,但在海量Topic、Partition场景会有性能问题
