# How RocketMQ-Rust Broker Stores and Manages Messages: Local File Store Architecture

> Discover how the rocketmq-broker stores and manages messages using its local file architecture. Learn about CommitLog, ConsumeQueue indexes, and key-based lookups for efficient message handling.

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

---

**The rocketmq-broker stores messages using a local-file architecture that writes sequentially to a CommitLog, builds logical ConsumeQueue indexes for fast topic-queue retrieval, and optionally maintains secondary indexes for key-based lookups.**

The rocketmq-rust broker (mxsm/rocketmq-rust) implements a high-performance message storage engine based on the classic RocketMQ design. It uses a **local-file store** as its default storage type, combining append-only physical logs with lightweight logical indexes to achieve both write throughput and efficient random reads. All storage components are orchestrated through the `BrokerRuntime` in [`rocketmq-broker/src/broker_runtime.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-broker/src/broker_runtime.rs).

## Storage Architecture Overview

The broker's storage subsystem centers on three core components that form a processing pipeline:

- **CommitLog**: An append-only sequential log file 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) that holds the raw binary representation of each `MessageExtBrokerInner`. Every write operation appends to this log first.
- **ConsumeQueue**: A logical index per topic-queue combination defined in [`rocketmq-store/src/queue/consume_queue.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/queue/consume_queue.rs). Each entry stores the physical offset, size, and tags-code, enabling consumers to locate message data without scanning the entire CommitLog.
- **IndexService**: An optional secondary index in [`rocketmq-store/src/index/index_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/index/index_service.rs) that supports fast lookups by message key or timestamp.

These structures are instantiated and wired together inside `BrokerRuntime::initialize_message_store`. When `store_type` equals `StoreType::LocalFile`, the runtime creates a `LocalFileMessageStore` and registers the dispatcher chain:

```rust
// In broker_runtime.rs – initialization of the local file store
if self.inner.message_store_config.store_type == StoreType::LocalFile {
    let mut message_store = ArcMut::new(LocalFileMessageStore::new(
        self.inner.message_store_config.clone(),
        self.inner.broker_config.clone(),
        self.inner.topic_config_manager().topic_config_table(),
        self.inner.broker_stats_manager.clone(),
        false,
    ));

    // Register dispatchers for ConsumeQueue and Index building
    let filter = Arc::new(CommitLogDispatcherCalcBitMap::new(
        self.inner.broker_config.clone(),
        self.inner.consumer_filter_manager.clone().unwrap(),
    ));
    self.inner.message_store_unchecked_mut().add_first_dispatcher(filter);
}

```

## The Message Write Pipeline

When a producer sends a message, the broker executes a multi-stage write pipeline that ensures durability while building the necessary indexes for retrieval.

### 1. Entry Point via SendMessageProcessor

The `SendMessageProcessor` in [`rocketmq-broker/src/processor/send_message_processor.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-broker/src/processor/send_message_processor.rs) receives the RPC request, constructs a `MessageExtBrokerInner`, and invokes `MessageStore::put_message`.

### 2. LocalFileMessageStore Persistence

The `LocalFileMessageStore::put_message` method 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) executes the core write logic:

```rust
async fn put_message(&mut self, mut msg: MessageExtBrokerInner) -> PutMessageResult {
    // Execute pre-put hooks (ACL, transaction checks)
    for hook in self.put_message_hook_list.iter() {
        if let Some(r) = hook.execute_before_put_message(&mut msg) { return r; }
    }

    // Append to CommitLog
    let result = self.commit_log.put_message(msg).await;

    // Trigger async index building via ReputMessageService
    if result.is_ok() {
        self.reput_message_service.notify_new_message();
    }
    result
}

```

### 3. CommitLog Sequential Write

The `CommitLog::put_message` implementation appends the serialized message to the current memory-mapped file and returns the physical offset and size.

### 4. Asynchronous Index Dispatch

The `ReputMessageService`, started within `LocalFileMessageStore::start`, reads the newly written segment, creates a `DispatchRequest`, and forwards it to the dispatcher chain (`CommitLogDispatcherDefault`). This chain contains two primary dispatchers:

- **`CommitLogDispatcherBuildConsumeQueue`**: Writes entries to the ConsumeQueue via `ConsumeQueueStore::append`.
- **`CommitLogDispatcherBuildIndex`**: Updates the `IndexService` for key-based queries.

## The Message Read Path

Consumer pull requests follow an efficient indexed retrieval pattern that avoids full CommitLog scans.

### 1. PullMessageProcessor Coordination

The `PullMessageProcessor` in [`rocketmq-broker/src/processor/pull_message_processor.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-broker/src/processor/pull_message_processor.rs) handles consumer requests. It locates the appropriate ConsumeQueue using `BrokerRuntime::find_consume_queue`.

### 2. ConsumeQueue Iteration

The processor iterates ConsumeQueue entries starting from the requested logical offset. Each entry contains the physical offset and size required to retrieve the actual message data.

### 3. CommitLog Random Read

For each valid ConsumeQueue entry passing the `MessageFilter`, the broker calls `CommitLog::get_message(phy_offset, size)` to fetch raw bytes. The `LocalFileMessageStore::get_message_with_size_limit` method orchestrates this flow:

```rust
// Core read logic from local_file_message_store.rs
let consume_queue = self.find_consume_queue(topic, queue_id);
if let Some(consume_queue) = consume_queue {
    let buffer_consume_queue = consume_queue.iterate_from_with_count(next_begin_offset, max_msg_nums);
    
    while let Some(cq_unit) = buffer_consume_queue.next() {
        let offset_py = cq_unit.pos;
        let size_py = cq_unit.size;

        // Apply tag filters
        if let Some(filter) = message_filter.as_ref() {
            if !filter.is_matched_by_consume_queue(...) { continue; }
        }

        // Physical retrieval from CommitLog
        if let Some(select_result) = self.commit_log.get_message(offset_py, size_py) {
            get_result.as_mut().unwrap().add_message(
                select_result,
                cq_unit.queue_offset as u64,
                cq_unit.batch_num as i32,
            );
        }
    }
}

```

## Advanced Storage Management Features

Beyond the core write-read cycle, the broker implements several specialized storage mechanisms:

### Timer Message Store

The `TimerMessageStore` in [`rocketmq-store/src/timer/timer_message_store.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/timer/timer_message_store.rs) manages delayed and scheduled messages. When a message with a future delivery timestamp arrives, it is written to a dedicated timer file rather than the standard CommitLog. A background thread in the timer wheel monitors ready messages and injects them into the normal `CommitLog` and `ConsumeQueue` pipeline when their delivery time arrives, at which point they become visible to consumers.

### High Availability (HA) Replication

Master-slave replication is implemented in [`rocketmq-store/src/ha/general_ha_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/general_ha_service.rs) and [`default_ha_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/default_ha_service.rs). The master embeds HA connection information in each CommitLog file header. The slave service establishes a TCP connection to the master, pulls binary data from the master's CommitLog, and writes it to its local store. The slave then runs the same dispatcher chain (`CommitLogDispatcherBuildConsumeQueue` and `CommitLogDispatcherBuildIndex`) to rebuild its local ConsumeQueue and indexes, ensuring data consistency across the replication group.

### Compaction Service

For topics configured with `CleanupPolicy::COMPACTION`, the `CompactionService` in [`rocketmq-store/src/kv/compaction_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/kv/compaction_service.rs) periodically rewrites storage files to discard obsolete message versions, retaining only the latest entry per key.

### Message Store Hooks and Stats

The broker supports custom validation logic through hooks defined in `rocketmq-broker/src/hook/`. The [`check_before_put_message.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/check_before_put_message.rs) hook enables ACL checks, while [`schedule_message_hook.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/schedule_message_hook.rs) handles batch validation before persistence. Performance metrics including I/O latency and disk-fall-behind are recorded by `BrokerStats` in [`rocketmq-store/src/stats/broker_stats.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/stats/broker_stats.rs) and refreshed via scheduled tasks in `BrokerRuntime::initialize_scheduled_tasks`.

## Code Examples

### Starting a Broker Programmatically

```rust
use rocketmq_common::common::broker::broker_config::BrokerConfig;
use rocketmq_store::config::message_store_config::MessageStoreConfig;
use rocketmq_broker::broker_runtime::BrokerRuntime;
use std::sync::Arc;

#[tokio::main]
async fn main() {
    // Load configuration (normally parsed from broker.toml)
    let broker_cfg = Arc::new(BrokerConfig::default());
    let store_cfg = Arc::new(MessageStoreConfig::default());

    // Create runtime with LocalFileMessageStore, CommitLog, and CQ
    let mut runtime = BrokerRuntime::new(broker_cfg.clone(), store_cfg.clone());

    // Initialize metadata and storage components
    runtime.initialize().await;

    // Start remoting server, HA services, and scheduled tasks
    runtime.start().await;
}

```

### Putting a Message (Producer Flow)

```rust
use rocketmq_common::common::message::message_ext_broker_inner::MessageExtBrokerInner;
use rocketmq_store::base::message_store::MessageStore;

async fn produce_message(runtime: &mut BrokerRuntime) {
    let mut msg = MessageExtBrokerInner::default();
    msg.topic = "TestTopic".into();
    msg.body = b"Hello RocketMQ".to_vec();
    
    if let Some(store) = runtime.inner.message_store.as_mut() {
        let result = store.put_message(msg).await;
        println!("Put result: {:?}", result);
    }
}

```

### Pulling Messages (Consumer Flow)

```rust
use rocketmq_common::common::message::message_ext::MessageExt;
use rocketmq_common::common::utils::cheetah_string::CheetahString;
use rocketmq_store::base::message_store::MessageStore;

async fn consume_messages(runtime: &BrokerRuntime) {
    let group = CheetahString::from("ConsumerGroupA");
    let topic = CheetahString::from("TestTopic");
    let queue_id = 0;
    let offset = 0;
    let max_msg = 32;

    if let Some(store) = runtime.inner.message_store.as_ref() {
        if let Some(result) = store
            .get_message(&group, &topic, queue_id, offset, max_msg, None)
            .await
        {
            for msg in result.msgs {
                println!("Payload: {:?}", msg.body);
            }
            println!("Next offset: {}", result.next_begin_offset);
        }
    }
}

```

## Summary

- **The rocketmq-broker uses a local-file store architecture** combining sequential CommitLog writes with logical ConsumeQueue indexes for optimal performance.
- **Message writes flow through three stages**: `SendMessageProcessor` → `LocalFileMessageStore::put_message` → `CommitLog`, followed by asynchronous dispatch to build ConsumeQueue and Index entries via `ReputMessageService`.
- **Message reads leverage ConsumeQueue indexes** to locate physical offsets in the CommitLog without full scans, as implemented in `LocalFileMessageStore::get_message_with_size_limit`.
- **Advanced features include** timer-based delayed messaging in `TimerMessageStore`, master-slave HA replication, compaction for `CleanupPolicy::COMPACTION` topics, and pre-put hooks for validation.
- **Core implementation files** reside 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), [`rocketmq-store/src/queue/consume_queue.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/queue/consume_queue.rs), and [`rocketmq-broker/src/broker_runtime.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-broker/src/broker_runtime.rs).

## Frequently Asked Questions

### How does the rocketmq-broker ensure message durability?

The broker achieves durability through the **CommitLog** append-only architecture 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). Every message is written sequentially to disk via memory-mapped files before any acknowledgment is sent to the producer. The `LocalFileMessageStore` synchronizes writes and maintains a `ReputMessageService` to ensure indexes are built only after the physical data is persisted.

### What is the difference between CommitLog and ConsumeQueue?

The **CommitLog** stores the actual message payload in a single sequential file per broker, optimizing write throughput. The **ConsumeQueue** in [`rocketmq-store/src/queue/consume_queue.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/queue/consume_queue.rs) is a lightweight logical index containing only fixed-size entries (physical offset, size, tags-code) organized per topic-queue combination. Consumers read the ConsumeQueue to find message locations, then retrieve full data from the CommitLog, enabling efficient random access without scanning the entire physical log.

### How does the broker handle delayed or scheduled messages?

Delayed messages are managed by the **TimerMessageStore** in [`rocketmq-store/src/timer/timer_message_store.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/timer/timer_message_store.rs). When a message with a future delivery timestamp arrives, it is written to a dedicated timer file rather than the standard CommitLog. A background thread in the timer wheel monitors ready messages and injects them into the normal `CommitLog` and `ConsumeQueue` pipeline when their delivery time arrives, at which point they become visible to consumers.

### Where is the HA (High Availability) replication logic implemented?

Master-slave replication is implemented in [`rocketmq-store/src/ha/general_ha_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/general_ha_service.rs) and [`default_ha_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/default_ha_service.rs). The master embeds HA connection information in each CommitLog file header. The slave service establishes a TCP connection to the master, pulls binary data from the master's CommitLog, and writes it to its local store. The slave then runs the same dispatcher chain (`CommitLogDispatcherBuildConsumeQueue` and `CommitLogDispatcherBuildIndex`) to rebuild its local ConsumeQueue and indexes, ensuring data consistency across the replication group.