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

> Explore how rocketmq-store implements local storage in RocketMQ-Rust using MappedFile append-only CommitLog and ConsumeQueue indexes. Learn about checkpoint recovery.

- Repository: [mxsm/rocketmq-rust](https://github.com/mxsm/rocketmq-rust)
- Tags: deep-dive
- Published: 2026-03-07

---

**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`](https://github.com/mxsm/rocketmq-rust/blob/main/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 |

```rust
// 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`](https://github.com/mxsm/rocketmq-rust/blob/main/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.

```rust
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`](https://github.com/mxsm/rocketmq-rust/blob/main/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:

```rust
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`](https://github.com/mxsm/rocketmq-rust/blob/main/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`](https://github.com/mxsm/rocketmq-rust/blob/main/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:

```rust
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`](https://github.com/mxsm/rocketmq-rust/blob/main/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`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/ha_service.rs)) replicates committed bytes to slave brokers before acknowledging producers.

```rust
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`](https://github.com/mxsm/rocketmq-rust/blob/main/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.

```rust
self.recover_normally(last_exit_ok).await;

```

## Practical Implementation Example

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

```rust
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`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/message_store/local_file_message_store.rs) | Central façade implementing `MessageStore` trait |
| [`rocketmq-store/src/log_file/commit_log.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/log_file/commit_log.rs) | Append-only commit log and HA/flush handling |
| [`rocketmq-store/src/log_file/mapped_file.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/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`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/queue/consume_queue_store.rs) | Management of per-queue index files |
| [`rocketmq-store/src/queue/consume_queue.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/queue/consume_queue.rs) | Low-level ConsumeQueue file format implementation |
| [`rocketmq-store/src/base/store_checkpoint.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/base/store_checkpoint.rs) | Persistent checkpoint for offsets and timestamps |
| [`rocketmq-store/src/base/flush_manager.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/base/flush_manager.rs) | Flush strategy implementation (sync/async) |
| [`rocketmq-store/src/ha/ha_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/ha_service.rs) | High-availability replication interface |
| [`rocketmq-store/src/config/store_path_config_helper.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/config/store_path_config_helper.rs) | Filesystem layout path helpers |
| [`rocketmq-store/src/utils/store_util.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/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.