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 regionselect_mapped_buffer()returns aSelectMappedBufferResultpointing 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:
- CommitLog::load() builds the MappedFileQueue (using optimized parallel loading when feature flags permit)
- CommitLog::recover_normally_optimized() scans commit-log files batch-wise, rebuilding ConsumeQueue and Index entries by replaying messages
- 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
MessageStoretrait for the broker. - Checkpoint-based recovery ensures consistency after restarts, using the
abortfile 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →