How rocketmq-store Implements Local Storage in RocketMQ-Rust: Architecture Deep Dive

The rocketmq-store module implements local storage using an append-only CommitLog backed by memory-mapped files (MappedFile), paired with ConsumeQueue index files for logical offsets, coordinated by LocalFileMessageStore with checkpoint-based recovery.

The rocketmq-store crate in the mxsm/rocketmq-rust repository provides a pure-Rust implementation of RocketMQ's durable message storage layer. It persists all message data to the local filesystem using a sophisticated architecture of memory-mapped commit logs, logical queue indexes, and transactional checkpoints. This implementation enables zero-copy I/O operations while ensuring strict durability guarantees through configurable flush strategies and high-availability replication.

Store Layout and Directory Structure

The local storage implementation organizes data under a root directory (store_path_root_dir) with a strict filesystem hierarchy. The StorePathConfigHelper utility in rocketmq-store/src/store_path_config_helper.rs defines this layout through functions like get_store_path_consume_queue(), which returns standardized subdirectories:

Folder Purpose
commitlog Raw binary commit-log files (e.g., 00000000000000000000) created by the CommitLog component
consumequeue Per-topic/queue logical index entries under <topic>/<queue_id> subdirectories
consumequeue_ext Extension files for batch ConsumeQueue operations
index Optional searchable index files for message key retrieval
checkpoint Persistent store checkpoint tracking maximum offsets and timestamps
abort Temporary file detecting unclean shutdowns
lock File-level lock preventing concurrent process access
// From rocketmq-store/src/store_path_config_helper.rs
pub fn get_store_path_consume_queue(root_dir: &str) -> String {
    PathBuf::from(root_dir).join("consumequeue")
        .to_string_lossy().into_owned()
}

Core Storage Architecture

LocalFileMessageStore: The Central Facade

The LocalFileMessageStore struct in rocketmq-store/src/message_store/local_file_message_store.rs acts as the primary facade, implementing the MessageStore trait. It coordinates all storage subsystems including the CommitLog, ConsumeQueueStore, IndexService, FlushManager, and HAService.

pub struct LocalFileMessageStore {
    message_store_config: Arc<MessageStoreConfig>,
    broker_config: Arc<BrokerConfig>,
    commit_log: ArcMut<CommitLog>,
    consume_queue_store: ArcMut<ConsumeQueueStore>,
    // … additional services
}

This structure provides the main API surface with methods like load(), start(), put_message(), and get_message(), wiring together the disparate storage components into a cohesive engine.

CommitLog: Append-Only Message Storage

The CommitLog component in rocketmq-store/src/log_file/commit_log.rs maintains an append-only byte-wise log of every message. It manages an ArcMut<MappedFileQueue> containing MappedFile objects, each representing a physical file sized according to mapped_file_size_commit_log.

When put_message() receives a MessageExtBrokerInner, it encodes the message and appends bytes to the tail MappedFile. If the current file reaches capacity, the AllocateMappedFileService creates a new file:

let result = mapped_file.as_ref().unwrap().append_message(
    &mut msg,
    self.append_message_callback.as_ref(),
    &put_message_context,
);

After successful append, CommitLog invokes handle_disk_flush_and_ha() to manage durability and replication based on the FlushDiskType configuration (sync or async).

MappedFile: Zero-Copy Memory Mapping

The MappedFile abstraction in rocketmq-store/src/log_file/mapped_file.rs wraps the memmap2 crate to provide zero-copy file access. It enables direct memory-mapped I/O without user-space buffering:

  • append_message() writes directly into the memory-mapped region
  • select_mapped_buffer() returns a SelectMappedBufferResult pointing to the underlying slice, avoiding data copies during reads

This implementation achieves high-throughput message storage by leveraging the operating system's virtual memory manager for file I/O.

ConsumeQueueStore and ConsumeQueue: Logical Indexing

While CommitLog stores messages sequentially, ConsumeQueueStore in rocketmq-store/src/queue/consume_queue_store.rs maintains logical indexes mapping queue offsets to physical file locations. Each topic-queue combination receives a dedicated ConsumeQueue file containing fixed-size entries (ConsumeQueueUnit) with fields for physical offset, message size, and tags code.

After CommitLog appends a message, assign_queue_offset() writes the index entry:

self.consume_queue_store.assign_queue_offset(msg);

This separation of concerns allows consumers to track logical progress independently while the CommitLog handles raw storage, enabling efficient message replay and seek operations.

Flush Management and High Availability

The storage layer guarantees durability through the FlushManager in rocketmq-store/src/base/flush_manager.rs, which implements sync and async flush strategies. When BrokerRole::SyncMaster is configured, the HAService (defined in rocketmq-store/src/ha/ha_service.rs) replicates committed bytes to slave brokers before acknowledging producers.

match (self.message_store_config.flush_disk_type, need_handle_ha) {
    (FlushDiskType::SyncFlush, true) => { 
        // Sync flush to disk and replicate to slaves
    }
    // … additional variants
}

Recovery and Checkpoint Mechanism

Startup recovery begins with loading the StoreCheckpoint from rocketmq-store/src/base/store_checkpoint.rs, which stores the last flushed offsets. The recovery process follows these steps:

  1. CommitLog::load() builds the MappedFileQueue (using optimized parallel loading when feature flags permit)
  2. CommitLog::recover_normally_optimized() scans commit-log files batch-wise, rebuilding ConsumeQueue and Index entries by replaying messages
  3. Truncation of dirty ConsumeQueue files beyond the last valid physical offset ensures consistency

An abort file created at startup and deleted only after clean shutdown enables detection of unclean exits, triggering more aggressive recovery procedures when present.

self.recover_normally(last_exit_ok).await;

Practical Implementation Example

The following example demonstrates the complete lifecycle of rocketmq-store local storage:

use rocketmq_store::{
    base::message_store::MessageStore,
    config::message_store_config::MessageStoreConfig,
    message_store::local_file_message_store::LocalFileMessageStore,
};
use rocketmq_common::common::message::message_ext_broker_inner::MessageExtBrokerInner;
use std::sync::Arc;
use tokio::runtime::Runtime;

fn main() {
    // 1. Configure the store
    let store_cfg = Arc::new(MessageStoreConfig::default());
    let broker_cfg = Arc::new(Default::default());

    // 2. Instantiate LocalFileMessageStore
    let mut store = LocalFileMessageStore::new(
        store_cfg.clone(),
        broker_cfg,
        Default::default(),
        None,
        false,
    );

    // 3. Initialize and start (async)
    let rt = Runtime::new().unwrap();
    rt.block_on(async {
        store.load().await;
        store.start().await.unwrap();
    });

    // 4. Produce a message
    let mut msg = MessageExtBrokerInner::new();
    msg.set_topic("TestTopic".into());
    msg.set_body(b"Hello RocketMQ".to_vec().into());

    rt.block_on(async {
        let put_result = store.put_message(msg).await;
        assert!(put_result.is_ok());
    });

    // 5. Consume a message
    rt.block_on(async {
        let get_res = store
            .get_message(&"GROUP".into(), &"TestTopic".into(), 0, 0, 32, None)
            .await
            .unwrap();
        println!("Fetched {} bytes", get_res.buffer_total_size());
    });
}

Key Source Files

File Role
rocketmq-store/src/message_store/local_file_message_store.rs Central façade implementing MessageStore trait
rocketmq-store/src/log_file/commit_log.rs Append-only commit log and HA/flush handling
rocketmq-store/src/log_file/mapped_file.rs Memory-mapped file abstraction for zero-copy I/O
rocketmq-store/src/queue/consume_queue_store.rs Management of per-queue index files
rocketmq-store/src/queue/consume_queue.rs Low-level ConsumeQueue file format implementation
rocketmq-store/src/base/store_checkpoint.rs Persistent checkpoint for offsets and timestamps
rocketmq-store/src/base/flush_manager.rs Flush strategy implementation (sync/async)
rocketmq-store/src/ha/ha_service.rs High-availability replication interface
rocketmq-store/src/config/store_path_config_helper.rs Filesystem layout path helpers
rocketmq-store/src/utils/store_util.rs System constants (physical memory sizing)

Summary

  • rocketmq-store implements local storage through a dual-layer architecture: the CommitLog for sequential message persistence and ConsumeQueue files for logical indexing.
  • Memory-mapped files (MappedFile) enable zero-copy I/O operations, eliminating user-space buffer copies during message production and consumption.
  • LocalFileMessageStore coordinates initialization, recovery, and runtime operations, implementing the MessageStore trait for the broker.
  • Checkpoint-based recovery ensures consistency after restarts, using the abort file to detect unclean shutdowns and replaying commit logs to rebuild indexes.
  • Configurable flush strategies (sync/async) and HAService replication provide durability and high-availability guarantees according to broker configuration.

Frequently Asked Questions

How does rocketmq-store ensure message durability?

The implementation ensures durability through the FlushManager, which persists memory-mapped data to disk using either synchronous or asynchronous strategies configured via FlushDiskType. When synchronous flush is enabled, the system calls msync on the MappedFile before returning success to producers. Additionally, the HAService replicates data to slave brokers in SyncMaster mode, ensuring durability across multiple nodes before acknowledging write operations.

What is the role of MappedFile in the storage implementation?

MappedFile provides the foundation for zero-copy storage by wrapping the memmap2 crate to map physical CommitLog files directly into process virtual memory. This allows append_message() to write directly to kernel-managed pages and select_mapped_buffer() to return references to message data without copying bytes through user-space buffers, significantly reducing CPU overhead and garbage collection pressure during high-throughput operations.

How does recovery work after an unclean shutdown?

Recovery begins by checking for the abort file in the store root directory. If present, the system executes recover_normally() or recover_normally_optimized() in CommitLog, scanning commit-log files batch-wise to rebuild ConsumeQueue indexes and IndexService entries by replaying messages. The process truncates ConsumeQueue files at the last valid physical offset recorded in the StoreCheckpoint, ensuring logical consistency between the commit log and indexes before accepting new messages.

What distinguishes CommitLog from ConsumeQueue in the storage architecture?

CommitLog serves as the single source of truth, storing all messages sequentially in append-only memory-mapped files regardless of topic or queue. ConsumeQueue acts as a secondary index maintaining per-topic/per-queue logical offsets that map to physical locations in the CommitLog. This separation enables efficient message retrieval by consumers tracking logical progress while allowing the CommitLog to optimize for sequential write throughput, following the classic RocketMQ storage pattern.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →