
RocketMQ作为阿里开源的消息中间件深受广大开发者的喜爱而这其中一个很重要原因就是它处理消息和拉取消息的速度非常快那么问题来了RocketMQ为什么这么快呢接下来我将从以下10个方面来探讨一下RocketMQ这么快的背后原因本文是基于RocketMQ 4.9.x版本讲解批量发送消息RocketMQ在发送消息的时候支持一次性批量发送多条消息如下代码所示public class Producer { public static void main(String[] args) throws Exception { //创建一个生产者指定生产者组为 sanyouProducer DefaultMQProducer producer new DefaultMQProducer(sanyouProducer); // 指定NameServer的地址 producer.setNamesrvAddr(192.168.200.143:9876); // 启动生产者 producer.start(); //用以及集合保存多个消息 ListMessage messages new ArrayList(); messages.add(new Message(sanyouTopic, 三友的java日记 0.getBytes())); messages.add(new Message(sanyouTopic, 三友的java日记 1.getBytes())); messages.add(new Message(sanyouTopic, 三友的java日记 2.getBytes())); // 发送消息并得到消息的发送结果然后打印 SendResult sendResult producer.send(messages); System.out.printf(%s%n, sendResult); // 关闭生产者 producer.shutdown(); } }通过批量发送消息减少了RocketMQ客户端与服务端也就是Broker之间的网络通信次数提高传输效率不过在使用批量消息的时候需要注意以下三点每条消息的Topic必须都得是一样的不支持延迟消息和事务消息不论是普通消息还是批量消息总大小默认不能超过4m消息压缩RocketMQ在发送消息的时候当发现消息的大小超过4k的时候就会对消息进行压缩这是因为如果消息过大会对网络带宽造成压力不过需要注意的是如果是批量消息的话就不会进行压缩如下所示压缩消息除了能够减少网络带宽造成压力之外还能够节省消息存储空间RocketMQ在往磁盘存消息的时候并不会去解压消息而是直接将压缩后的消息存到磁盘消费者拉取到的消息其实也是压缩后的消息不过消费者在拿到消息之后会对消息进行解压缩当我们的业务系统拿到消息的时候其实就是解压缩后的消息虽然压缩消息能够减少带宽压力和磁盘存储压力但是由于压缩和解压缩的过程都是在客户端生产者、消费者完成的所以就会导致客户端消耗更多的CPU资源对CPU造成一定的压力高性能网络通信模型当生产者处理好消息之后就会将消息通过网络通信发送给服务端而RocketMQ之所以快的一个非常重要原因就是它拥有高性能网络通信模型RocketMQ网络通信这块底层是基于Netty来实现的Netty是一款非常强大、非常优秀的网络应用程序框架主要有以下几个优点异步和事件驱动Netty基于事件驱动的架构使用了异步I/O操作避免了阻塞式I/O调用的缺陷能够更有效地利用系统资源提高并发处理能力。高性能Netty针对性能进行了优化比如使用直接内存进行缓冲减少垃圾回收的压力和内存拷贝的开销提供了高吞吐量、低延迟的网络通讯能力。可扩展性Netty的设计允许用户自定义各种Handler来处理协议编码、协议解码和业务逻辑等。并且它的模块可插拔性设计使得用户可以根据需要轻松地添加或更换组件。简化API与Java原生NIO库相比Netty提供了更加简洁易用的API大大降低了网络编程的复杂度。安全Netty内置了对SSL/TLS协议的支持使得构建安全通信应用变得容易。丰富的协议支持Netty提供了HTTP、HTTP/2、WebSocket、Google Protocol Buffers等多种协议的编解码支持满足不同网络应用需求。...就是因为Netty如此的强大所以不仅仅RocketMQ是基于Netty实现网络通信的几乎绝大多数只要涉及到网络通信的Java类框架底层都离不开Netty的身影比如知名RPC框架Dubbo、Java gRPC实现、Redis的亲儿子Redisson、分布式任务调度平台xxl-job等等它们底层在实现网络通信时都是基于Netty框架零拷贝技术当消息达到RocketMQ服务端之后为了能够保证服务端重启之后消息也不丢失此时就需要将消息持久化到磁盘由于涉及到消息持久化操作就涉及到磁盘文件的读写操作RocketMQ为了保证磁盘文件的高性能读写使用到了一个叫零拷贝的技术1、传统IO读写方式说零拷贝之前先说一下传统的IO读写方式。比如现在有一个需求将磁盘文件通过网络传输出去那么整个传统的IO读写模型如下图所示传统的IO读写其实就是read write的操作整个过程会分为如下几步用户调用read()方法开始读取数据此时发生一次上下文从用户态到内核态的切换也就是图示的切换1将磁盘数据通过DMA拷贝到内核缓存区将内核缓存区的数据拷贝到用户缓冲区这样用户也就是我们写的代码就能拿到文件的数据read()方法返回此时就会从内核态切换到用户态也就是图示的切换2当我们拿到数据之后就可以调用write()方法此时上下文会从用户态切换到内核态即图示切换3CPU将用户缓冲区的数据拷贝到Socket缓冲区将Socket缓冲区数据拷贝至网卡write()方法返回上下文重新从内核态切换到用户态即图示切换4整个过程发生了4次上下文切换和4次数据的拷贝这在高并发场景下肯定会严重影响读写性能。所以为了减少上下文切换次数和数据拷贝次数就引入了零拷贝技术。2、零拷贝零拷贝技术是一个思想指的是指计算机执行操作时CPU不需要先将数据从某处内存复制到另一个特定区域。实现零拷贝的有以下两种方式mmap()sendfile()mmap()mmapmemory map是一种内存映射文件的方法即将一个文件或者其它对象映射到进程的地址空间实现文件磁盘地址和进程虚拟地址空间中一段虚拟地址的一一对映关系。简单地说就是内核缓冲区和应用缓冲区进行映射用户在操作应用缓冲区时就好像在操作内核缓冲区比如你往应用缓冲区写数据就好像直接往内核缓冲区写数据这个过程不涉及到CPU拷贝而传统IO就需要将在写完应用缓冲区之后需要将数据通过CPU拷贝到内核缓冲区同样地上述文件传输功能如果使用mmap的话由于我们可以直接操作内核缓冲区此时我们就可以将内核缓冲区的数据直接CPU拷贝到Socket缓冲区整个IO模型就会如下图所示基于mmap IO读写其实就变成mmap write的操作也就是用mmap替代传统IO中的read操作当用户发起mmap调用的时候会发生上下文切换1进行内存映射然后数据被拷贝到内核缓冲区mmap返回发生上下文切换2随后用户调用write发生上下文切换3将内核缓冲区的数据拷贝到Socket缓冲区write返回发生上下文切换4。上下文切换的次数仍然是4次但是拷贝次数只有3次少了一次CPU拷贝。所以总的来说使用mmap就可以直接少一次CPU拷贝。说了这么多那么在Java中如何去实现mmap也就是内核缓冲区和应用缓冲区映射呢其实在Java NIO类库中就提供了相应的API当然底层也还是调用Linux系统的mmap()实现的代码如下所示FileChannel fileChannel new RandomAccessFile(test.txt, rw).getChannel(); MappedByteBuffer mappedByteBuffer fileChannel.map(FileChannel.MapMode.READ_WRITE, 0, fileChannel.size());MappedByteBuffer你可以认为操作这个对象就好像直接操作内核缓冲区比如可以通过MappedByteBuffer读写磁盘文件此时就好像直接从内核缓冲区读写数据当然也可以直接通过MappedByteBuffer将文件的数据拷贝到Socket缓冲区实现上述文件传输的模型这里我就不贴相应的代码了RocketMQ在存储文件时就是通过mmap技术来实现高效的文件读写RocketMQ中使用mmap代码虽然前面一直说mmap不涉及CPU拷贝但在某些特定场景下尤其是在写操作或特定的系统优化策略下还是可能涉及CPU拷贝。sendfile()sendfile()跟mmap()一样也会减少一次CPU拷贝但是它同时也会减少两次上下文切换。sendfile()主要是用于文件传输比如将文件传输到另一个文件又或者是网络当基于sendfile()时一次文件传输的过程就如下图所示用户发起sendfile()调用时会发生切换1之后数据通过DMA拷贝到内核缓冲区之后再将内核缓冲区的数据CPU拷贝到Socket缓冲区最后拷贝到网卡sendfile()返回发生切换2。同样地Java NIO类库中也提供了相应的API实现sendfile当然底层还是操作系统的sendfile()FileChannel channel FileChannel.open(Paths.get(./test.txt), StandardOpenOption.WRITE, StandardOpenOption.CREATE); //调用transferTo方法向目标数据传输 channel.transferTo(position, len, target);FileChannel的transferTo方法底层就是基于sendfile来的在如上代码中并没有文件的读写操作而是直接将文件的数据传输到target目标缓冲区也就是说sendfile传输文件时是无法知道文件的具体的数据的但是mmap不一样mmap可以来直接修改内核缓冲区的数据假设如果需要对文件的内容进行修改之后再传输mmap可以满足小总结在传统IO中如果想将用户缓存区的数据放到内核缓冲区需要经过CPU拷贝而基于零拷贝技术可以减少CPU拷贝次数常见的有两种mmap()sendfile()mmap()是将用户缓冲区和内核缓冲区共享操作用户缓冲区就好像直接操作内核缓冲区读写数据时不需要CPU拷贝Java中可以使用MappedByteBuffer这个API来达到操作内核缓冲区的效果sendfile()主要是用于文件传输可以通过sendfile()将一个文件内容传输到另一个文件中或者是网络中sendfile()在整个过程中是无法对文件内容进行修改的如果想修改之后再传输可以通过mmap来修改内容之后再传输上面出现的API都是Java NIO标准类库中的如果你看的还是很迷糊那直接记住一个结论之所以基于零拷贝技术能够高效的实现文件的读写操作主要因为是减少了CPU拷贝次数和上下文切换次数在RocketMQ中底层是基于mmap()来实现文件的高效读写的顺序写RocketMQ在存储消息时除了使用零拷贝技术来实现文件的高效读写之外还使用顺序写的方式提高数据写入的速度RocketMQ会将消息按照顺序一条一条地写入文件中这种顺序写的方式由于减少了磁头的移动和寻道时间在大规模数据写入的场景下使得数据写入的速度更快高效的数据存储结构Topic和队列的关系在RocketMQ中默认会为每个Topic在每个服务端Broker实例上创建4个队列如果有两个Broker那么默认就会有8个队列每个Broker上的队列上的编号queueId都是从0开始CommitLog前面一直说当消息到达RocektMQ服务端时需要将消息存到磁盘文件RocketMQ给这个存消息的文件起了一个高大上的名字CommitLog由于消息会很多所以为了防止文件过大CommitLog在物理磁盘文件上被分为多个磁盘文件每个文件默认的固定大小是1G消息在写入到文件时除了包含消息本身的内容数据也还会包含其它信息比如消息的Topic消息所在队列的id生产者发送消息时会携带这个队列id消息生产者的ip和端口...这些数据会和消息本身按照一定的顺序同时写到CommitLog文件中上图中黄色排列顺序和实际的存的内容并非实际情况我只是举个例子ConsumeQueue除了CommitLog文件之外RocketMQ还会为每个队列创建一个磁盘文件RocketMQ给这个文件也起了一个高大上的名字ConsumeQueue当消息被存到CommitLog之后其实还会往这条消息所在队列的ConsumeQueue文件中插一条数据每个队列的ConsumeQueue也是由多个文件组成每个文件默认是存30万条数据插入ConsumeQueue中的每条数据由20个字节组成包含3部分信息消息在CommitLog的起始位置8个字节也被称为偏移量消息在CommitLog存储的长度4个字节消息tag的hashCode8个字节每条数据也有自己的编号offset默认从0开始依次递增所以通过ConsumeQueue中存的数据可以从CommitLog中找到对应的消息那么这个ConsumeQueue有什么作用呢其实通过名字也能猜到这其实跟消息消费有关当消费者拉取消息的时候会告诉服务端四个比较重要的信息自己需要拉取哪个Topic的消息从Topic中的哪个队列queueId拉取从队列的哪个位置offset拉取消息拉取多少条消息(默认32条)服务端接收到消息之后总共分为四步处理首先会找到对应的Topic之后根据queueId找到对应的ConsumeQueue文件然后根据offset位置从ConsumeQueue中读取跟拉取消息条数一样条数的数据由于ConsumeQueue每条数据都是20个字节所以根据offset的位置可以很快定位到应该从文件的哪个位置开始读取数据最后解析每条数据根据偏移量和消息的长度到CommitLog文件查找真正的消息内容整个过程如下图所示所以从这可以看出当消费者在拉取消息时ConsumeQueue其实就相当于是一个索引文件方便快速查找在CommitLog中的消息并且无论CommitLog存多少消息整个查找消息的时间复杂度都是O(1)由于ConsumeQueue每条数据都是20个字节所以如果需要找第n条数据只需要从第n * 20个字节的位置开始读20个字节的数据即可这个过程是O(1)的当从ConsumeQueue找到数据之后解析出消息在CommitLog存储的起始位置和大小之后就直接根据这两个信息就可以从CommitLog中找到这条消息了这个过程也是O(1)的所以整个查找消息的过程就是O(1)的所以从这就可以看出ConsumeQueue和CommitLog相互配合就能保证快速查找到消息消费者从而就可以快速拉取消息异步处理RocketMQ在处理消息时有很多异步操作这里我举两个例子异步刷盘异步主从复制异步刷盘前面说到文件的内容都是先写到内核缓冲区也可以说是PageCache而写到PageCache并不能保证消息一定不丢失因为如果服务器挂了这部分数据还是可能会丢失的所以为了解决这个问题RocketMQ会开启一个后台线程这个后台线程默认每隔0.5s会将消息从PageCache刷到磁盘中这样就能保证消息真正的持久化到磁盘中异步主从复制在RocketMQ中支持主从复制的集群模式这种模式下写消息都是写入到主节点读消息一般也是从主节点读但是有些情况下可能会从从节点读从节点在启动的时候会跟主节点建立网络连接当主节点将消息存储的CommitLog文件之后会通过后台一个异步线程不停地将消息发送给从节点从节点接收到消息之后就直接将消息存到CommitLog文件小总结就是因为有这些异步操作大大提高了消息存储的效率不过值得注意的尽管异步可以提高效率但是也增加了不确定性比如丢消息等等当然RocketMQ也支持同步等待消息刷盘和主从复制成功但这肯定会导致性能降低所以在项目中可以根据自己的业务需要选择对应的刷盘和主从复制的策略批量处理除了异步之外RocketMQ还大量使用了批量处理机制比如前面说过消费者拉取消息的时候可以指定拉取拉取消息的条数批量拉取消息这种批量拉取机制可以减少消费者跟RocketMQ服务端的网络通信次数提高效率除了批量拉取消息之外RocketMQ在提交消费进度的时候也使用了批量处理机制所谓的提交消费进度就是指当消费者在成功消费消息之后需要将所消费消息的offsetConsumeQueue中的offset提交给RocketMQ服务端告诉RocketMQ这个Queue的消息我已经消费到了这个位置了这样一旦消费者重启了或者其它啥的要从这个Queue重新开始拉取消息的时候此时他只需要问问RocketMQ服务端上次这个Queue消息消费到哪个位置了之后消费者只需要从这个位置开始消费消息就行了这样就解决了接着消费的问题RocketMQ在提交消费进度的时候并不是说每消费一条消息就提交一下这条消息对应的offset而是默认每隔5s定时去批量提交一次这5s钟消费消息的offset锁优化由于RocketMQ内部采用了很多线程异步处理机制这就一定会产生并发情况下的线程安全问题在这种情况下RocketMQ进行了多方面的锁优化以提高性能和并发能力就比如拿消息存储来说为了保证消息是按照顺序一条一条地写入到CommitLog文件中就需要对这个写消息的操作进行加锁而RocketMQ默认使用ReentrantLock来加锁并不是synchronized当然除了默认情况外RocketMQ还提供了一种基于CAS加锁的实现这种实现可以在写消息压力较低的情况下使用当然除了写消息之外在一些其它的地方RocketMQ也使用了基于CAS的原子操作来代替传统的锁机制例如使用大量使用了AtomicInteger、AtomicLong等原子类来实现并发控制避免了显式的锁竞争提高了性能线程池隔离RocketMQ在处理请求的时候会为不同的请求分配不同的线程池进行处理比如对于消息存储请求和拉取消息请求来说Broker会有专门为它们分配两个不同的线程池去分别处理这些请求这种让不同的业务由不同的线程池去处理的方式能够有效地隔离不同业务逻辑之间的线程资源的影响比如消息存储请求处理过慢并不会影响处理拉取消息请求所以RocketMQ通过线程隔离及时可以有效地提高系统的并发性能和稳定性