Java 消息队列持久化设计:从文件存储到数据可靠性保障
在实现消息队列时消息持久化是保证数据可靠性的核心环节 —— 即使服务重启、机器宕机已经投递的消息也不会丢失。本文将详细介绍我们自研 Java 消息队列中消息持久化的文件存储设计与实现思路。一、整体存储结构设计由于消息是依附于队列存在的我们可以按队列维度来组织持久化文件我们可以在data这个目录中再创建一些子目录并约定每个队列都有一个子目录子目录的名字就是队列的名字。每个队列的子目录之下再分配两个文件来存储信息第一个文件:queue_data.txt这里保存消息的内容.第二个文件:queue_stat.txt这里保存消息的统计信息.二、queue_data 文件格式设计这个文件中包含了若干个消息每个消息都以二进制的方式进行存储每个消息的前面的四个字节用来表示该条消息的实际长度接下来的一部分就是消息的二进制数据。我们的Message类实现了Serializable接口其核心结构如下public class Message implements Serializable { // 核心属性消息基本属性 消息体 private BasicProperties basicProperties new BasicProperties(); private byte[] body; // 辅助属性用于定位消息在文件中的位置不序列化 private transient long offsetBeg 0; // 消息开头距离文件头的偏移量 private transient long offsetEnd 0; // 消息结尾距离文件头的偏移量 // 逻辑删除标记0x1 表示有效0x0 表示已删除 private byte isValid 0x1; }其中BasicProperties类封装了消息的核心元数据public class BasicProperties implements Serializable { // 消息唯一 ID使用 UUID 保证唯一性 private String messageId; // 路由键用于交换机与队列的绑定匹配 private String routingKey; // 持久化模式1 表示不持久化2 表示持久化参考 RabbitMQ 设计 private int deliveryMode 1; }因此单条消息的二进制数据在文件中实际布局为isValid这个属性是用来标识当前这个消息是否在文件中是有效的。消息定位与逻辑删除偏移量定位offsetBeg和offsetEnd标记了消息在文件中的起止位置这两个属性被transient修饰不会被序列化到磁盘仅在内存中用于快速定位消息对应的磁盘位置。逻辑删除为了避免频繁修改文件带来的性能开销我们采用逻辑删除而非物理删除通过修改isValid字段为0x0标记消息为无效后续通过垃圾回收机制统一清理无效数据。三、垃圾回收GC机制随着消息不断写入和删除queue_data.txt会逐渐膨胀且其中会存在大量无效消息。为了回收磁盘空间我们需要实现垃圾回收机制1. GC 触发条件当文件中无效消息占比超过阈值如 50%时触发 GC。也可以按固定时间周期如每 2000 条消息写入后触发。2. GC 实现流程遍历读取遍历queue_data.txt读取所有标记为isValid 0x1的有效消息。数据重写将有效消息依次写入一个临时文件如queue_data.txt.tmp。文件替换删除原queue_data.txt将临时文件重命名为queue_data.txt。更新偏移量重新计算所有有效消息在新文件中的offsetBeg和offsetEnd更新内存中的Message对象。这种方式虽然会有一次全量数据拷贝但能最大程度保证文件紧凑避免磁盘空间浪费。四、queue_stat 统计文件设计queue_stat.txt是一个纯文本文件用于记录队列的统计信息格式为一行两列第一列是queue_data.txt中总的消息的数目第二列是queue_data.txt 中有效消息的数目两者使用\t分割形如:2000\t1500五、序列化方案选择在将Message对象转换为二进制时我们选择了Java 原生序列化ObjectOutputStream/ObjectInputStream// 把一个对象序列化成一个字节数组 public static byte[] toBytes(Object object) throws IOException{ // 1. 创建ByteArrayOutputStream对象可以把它理解成一个可自动扩容的字节桶 // 普通的byte[]数组长度是固定的没法边写边扩容而这个桶可以 // 序列化出来的二进制数据会先临时存在这个桶里最后统一转成byte[] try(ByteArrayOutputStream byteArrayOutputStreamnew ByteArrayOutputStream()){ // 2. 创建ObjectOutputStream对象这个是序列化工具 // 它的作用就是把Java对象转换成二进制字节数据 // 并且把转换后的数据写入到上面的字节桶里 try(ObjectOutputStream objectOutputStreamnew ObjectOutputStream(byteArrayOutputStream)){ // 3. 核心操作调用writeObject把对象序列化 // 这一步会把传入的object拆成二进制字节然后写入到ObjectOutputStream // 因为ObjectOutputStream和前面的字节桶关联着所以最终数据会落到桶里 objectOutputStream.writeObject(object); } // 4. 把字节桶里存的所有二进制数据转成普通的byte[]数组返回 // 到这一步对象就已经完全变成可以存文件/传网络的字节数组了 return byteArrayOutputStream.toByteArray(); } } // 把一个字节数组, 反序列化成一个对象 public static Object fromBytes(byte[] data) throws IOException, ClassNotFoundException { // 定义要返回的对象 Object objectnull; // 1. 创建ByteArrayInputStream对象你可以把它理解成一个字节读取器 // 它会把传入的byte[]数组包装起来提供按顺序读取字节的能力 try(ByteArrayInputStream byteArrayInputStreamnew ByteArrayInputStream(data)){ // 2. 创建ObjectInputStream对象这个是反序列化工具 // 它的作用就是把二进制字节数据还原成Java对象 // 它会从上面的字节读取器里拿字节数据来解析 try(ObjectInputStream objectInputStreamnew ObjectInputStream(byteArrayInputStream)){ // 3. 核心操作调用readObject把字节数组反序列化 // 这一步会把二进制字节重新组装成原来的Java对象赋值给object变量 objectobjectInputStream.readObject(); } } // 4. 返回还原后的对象 return object; }结语本文围绕自研 Java 消息队列的消息持久化需求完整阐述了基于文件存储的设计与实现方案通过按队列维度拆分存储目录明确 queue_data 与 queue_stat 文件的职责边界借助二进制存储 逻辑删除的方式平衡读写性能与空间利用率通过垃圾回收机制解决无效数据占用磁盘的问题同时选择 Java 原生序列化完成对象与二进制的转换保障数据落地的完整性。