ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

RocketMQ源码级实现原理02-存储系统概要设计

RocketMQ源码级实现原理02-存储系统概要设计 最新的NVME协议的SSD硬盘磁盘存储文件写入速度顺序读写 3g/s左右随机读写2g/s左右内存条的读写速度是10个g左右存储系统概要设计主要就是讲消息从producer到达broker后怎么落地存储CommitLog的写是随机写但是读可能会有随机读当然也可能直接读pagecache大方向是对的——这是两者存储模型根本差异决定的经典对比。但RocketMQ 完全不会下降要打个折扣。我按代码和架构说清楚。Kafka 为什么会随 topic更准确说是分区数下降Kafka 的存储模型是每个 topic-partition 是一个独立的日志目录 独立的一组文件.log / .index / .timeindex。- 写入时broker 同时在给多个分区的文件追加 → 分区数一多写入就从顺序写一个大文件变成在几十上百个文件之间来回写即随机 I/O- 每个分区都要占用文件句柄和 page cache分区多了缓存被切碎、命中率下降且各分区的尾部数据互相挤占缓存- 每个分区还有独立的索引、元数据分区数增长还会带来更长的 rebalance、leader 选举、controller 元数据压力等这些和写入无直接关系但拖累整体。所以 Kafka 在分区数不是单纯的 topic 数到几百上千时写入性能会明显下滑——社区普遍建议单 broker 分区数别太多。到 100 这个量级是否明显取决于磁盘HDD明显、SSD 弱化、缓存大小和配置。RocketMQ 为什么写入路径几乎不受 topic 数影响RocketMQ 用的是**单 CommitLog 消费队列索引**模型1. 所有 topic 的所有消息全部追加到同一个 CommitLog 文件DefaultMessageStore.commitLog new CommitLog(this)putMessage就是往这一个文件顺序追加——写入永远是单文件顺序写跟 topic 数量无关2. 追加完成后由一个后台线程ReputMessageServiceDefaultMessageStore.java:1864从 commitlog 顺序读取再把消息位置分发写进各个 topic-queue 对应的ConsumeQueue3. ConsumeQueue 每条记录是定长 20 字节ConsumeQueue.CQ_STORE_UNIT_SIZE 20只是物理偏移大小tag hash很小。关键点真正贵的那一步消息数据落盘始终是单文件顺序写不管你有 1 个还是 1000 个 topic。多 topic 只影响第二步分发索引而那是小的定长记录、且是异步的。这就是RocketMQ 敢宣称支持海量 topic的底气。一个好的文件系统它的性能肯定是要优于分布式KV和newSQL;Kafka和RocketMQ存储模型的差别Kafka会有多个文件存储不像rocketmq总是写一个文件能实现一直往文件末尾追加的随机写RocketMQ存储架构图1. ConsumeQueue里面放的是一个个索引条目索引条目有三个字段CommitLogOffSet MessageSize、TagHashCode2. doDispatch 异步线程 构建ConsumeQueueCommitLog底层结构这么多MappedFile如何管理起来对外暴露就用一个MappedFileQueueMessage格式存储模块实现总体代码层级设计架构存储这块的设计很好层次很清晰存储逻辑层比如CommitLog类中出了对外暴露putMessage()的接口外还需要提供一些内部类比如FlushRealTimeService这类的刷盘线程逻辑存储IO层则直接和PageCache、和磁盘打交道各业务拥有独立线程池也就是不同的逻辑处理RocketMQ都会为它分配不同的线程池比如这里为了处理发送端发过来的消息写入磁盘就通过sendMssageExecutor线程池来执行sendMssageProcessor中的代码完成一条消息的写入磁盘CommitLog写入模块未关闭自动创建topic开关生产中需要关闭自动创建topicDefaultMessageStore#putMessage()可以看到这里指定了每个mappedFile大小是10M实际中这个大小是1GcommitLog作为共享资源写入要加锁commitLog.putMessage()会有加锁释放锁的逻辑因为要保证同一时间只有一个线程去往MappedFile中写入消息数据byteBuffer.position()代表的就是当前在00000000000000000174080这个commitlog文件中的相对偏移表示在当前要写入的消息之前上一条消息已经写到了174080 byteBuffer.position()的位置了所以当前要写入的消息是从174080 byteBuffer.position()往后开始写入因为一个消息不能跨两个mappedFile所以当最后一个mappedFile的剩余空间小于 当前文件大小 8个字节的结尾魔数则就要开始创建下一个mappedFile把当前消息写入下一个mappedFile了ConsumeQueue写入模块物理结构Broker端每个topic下的各个队列queue持久化后的消费进度数据主要代码组件有独立的线程FlushConsumeQueueService线程把构建好的consumeQueue对应的mappedFile刷入磁盘代码实现已经刷过盘的消息offset - 当前已经reput过的消息offset就是剩余的还需要reput的消息
返回列表