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

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.

Storage Architecture Overview

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

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:

// 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 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 executes the core write logic:

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 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:

// 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 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 and 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 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 hook enables ACL checks, while 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 and refreshed via scheduled tasks in BrokerRuntime::initialize_scheduled_tasks.

Code Examples

Starting a Broker Programmatically

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)

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)

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, rocketmq-store/src/queue/consume_queue.rs, and 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. 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 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. 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 and 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.

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 →