Data Durability Features in RocketMQ-Rust: Write-Ahead Logging, Flush Strategies, and HA Replication

RocketMQ-Rust guarantees message durability through a write-ahead CommitLog architecture with configurable sync/async flush strategies, group-commit batching for reduced syscall overhead, periodic checkpointing for crash recovery, and High-Availability (HA) replication with configurable In-Sync Replica (ISR) semantics.

The mxsm/rocketmq-rust repository implements enterprise-grade message persistence by combining memory-mapped sequential I/O with tunable durability guarantees. Understanding these data durability features in RocketMQ-Rust is essential for operators configuring brokers for either low-latency streaming or strong-consistency workloads.

Write-Ahead CommitLog Architecture

All inbound messages are appended to the CommitLog, a sequential, memory-mapped file that stores raw message bytes. The CommitLog implementation resides 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).

During a putMessage operation, the log entry is written to the mapped file and subsequently flushed according to the configured flush mode. This write-ahead approach ensures that data exists on disk before acknowledging producers, forming the foundation of the durability guarantees.

Flush Disk Types and Configuration

The MessageStoreConfig struct defined in [rocketmq-store/src/config/message_store_config.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/config/message_store_config.rs) controls the global flush policy via the flush_disk_type field. This field accepts the FlushDiskType enum from [rocketmq-store/src/config/flush_disk_type.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/config/flush_disk_type.rs).

SyncFlush Mode

When configured for SyncFlush, the writer blocks until the operating system flushes the file to persistent storage. This mode is implemented in DefaultFlushManager::handle_disk_flush within [rocketmq-store/src/log_file/flush_manager_impl/default_flush_manager.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/log_file/flush_manager_impl/default_flush_manager.rs).

SyncFlush provides the strongest durability guarantee, ensuring that acknowledged messages survive process crashes or power failures, though at the cost of higher write latency.

AsyncFlush Mode

AsyncFlush allows the writer to return immediately after writing to the memory-mapped buffer. A background thread periodically invokes MappedFile::flush to persist data to disk. This implementation shares the same location in DefaultFlushManager but follows the asynchronous execution branch.

AsyncFlush maximizes throughput and minimizes latency, though it risks losing unflushed messages in the event of an unclean shutdown. This mode suits scenarios where raw throughput outweighs strict durability requirements.

Group-Commit Service for Synchronous Flush

To optimize the synchronous path, RocketMQ-Rust employs a GroupCommitService that batches flush requests. Rather than triggering an fsync for every individual message, the service accumulates requests and wakes the flush thread once per batch, significantly reducing syscall overhead.

The service is instantiated in DefaultFlushManager::new and processes GroupCommitRequest objects. This batching mechanism makes SyncFlush viable for high-throughput scenarios without sacrificing the durability guarantees of synchronous disk writes.

Checkpointing and Recovery

The system periodically writes a store checkpoint that records the highest committed offsets for the CommitLog, ConsumeQueue, and other index structures. The StoreCheckpoint struct is defined in [rocketmq-store/src/base/store_checkpoint.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/base/store_checkpoint.rs).

Checkpoint flushing integrates with the flush_manager pipeline, ensuring that recovery after a crash can restore the broker to a consistent state. The checkpoint references the last known good positions, allowing the system to truncate partially written entries and resume from valid offsets.

High-Availability Replication and ISR

Beyond local disk persistence, RocketMQ-Rust implements High-Availability (HA) replication to guard against node failures. When operating as a SyncMaster, the broker replicates messages to slave nodes before acknowledging producers.

HA Service Architecture

The HA flow triggers inside CommitLog::handle_disk_flush_and_ha via the need_handle_ha check. Replication proceeds by sending GroupCommitRequest objects to the HA service (ha_service.put_request), defined in [rocketmq-store/src/ha/ha_service.rs](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/ha_service.rs). Concrete implementations reside in rocketmq-store/src/ha/default_ha_service.rs and related modules.

In-Sync Replica Management

Durability across the cluster depends on In-Sync Replica (ISR) semantics. The required acknowledgment count derives from message_store_config.in_sync_replicas, with optional auto-adjustment based on current replica lag via CommitLog::calc_need_ack_nums (lines 78-86 in the source).

A write receives confirmation only after the configured ISR count persists the data, ensuring quorum-based durability across distributed nodes.

Transactional Durability Guarantees

For transactional message flows, the CommitLog updates the ConsumeQueue only after the transaction commit entry is flushed to disk. This prevents partial visibility of uncommitted messages. The logic resides in increase_offset and assign_offset methods within the CommitLog implementation, ensuring atomicity between transaction state and message visibility.

Practical Configuration Examples

Configure Synchronous Flushing with Quorum Replication

To enforce the strongest durability guarantees—synchronous disk flush and acknowledgment from two replicas:

use rocketmq_store::config::message_store_config::MessageStoreConfig;
use rocketmq_store::config::flush_disk_type::FlushDiskType;

let mut cfg = MessageStoreConfig::default();
cfg.flush_disk_type = FlushDiskType::SyncFlush;    // block until OS flushes
cfg.in_sync_replicas = 2;                         // require 2 replicas ACK
cfg.enable_auto_in_sync_replicas = false;          // use static ISR size
// … start the message store with `cfg`

Enable Asynchronous Flushing for Lower Latency

For scenarios prioritizing throughput over crash durability:

let mut cfg = MessageStoreConfig::default();
cfg.flush_disk_type = FlushDiskType::AsyncFlush; // fire‑and‑forget flush
cfg.flush_commit_log_timed = true;               // periodic flush thread
// start the broker – messages are considered durable once HA (if enabled) confirms

Disable HA for Stand-Alone Deployment

use rocketmq_common::base::broker_role::BrokerRole;

let mut cfg = MessageStoreConfig::default();
cfg.broker_role = BrokerRole::AsyncMaster; // no HA
cfg.duplication_enable = false; // disables HA checks in `handle_disk_flush_and_ha`

Inspect Current Checkpoint Offsets

To verify the last persisted offset programmatically:

let checkpoint = message_store.get_store_checkpoint(); // implements trait in base/message_store.rs
println!("Last flushed offset: {}", checkpoint.get_commit_log_offset());

Summary

  • Write-Ahead CommitLog: All messages append to a memory-mapped sequential log in commit_log.rs before acknowledgment.
  • Configurable Flush: Choose between SyncFlush (blocking, durable) and AsyncFlush (background, fast) via MessageStoreConfig.
  • Group-Commit Optimization: The GroupCommitService batches sync flush requests to reduce syscall overhead.
  • Checkpoint Recovery: StoreCheckpoint persists committed offsets enabling crash recovery to consistent states.
  • HA and ISR: Replication through HAService with configurable in_sync_replicas ensures distributed durability.
  • Transaction Safety: ConsumeQueue updates occur only after transaction commit entries are flushed.

Frequently Asked Questions

What is the difference between SyncFlush and AsyncFlush in RocketMQ-Rust?

SyncFlush blocks the producer thread until the operating system confirms the data is on persistent storage, providing the strongest durability guarantee against crashes. AsyncFlush returns immediately to the producer, with a background thread handling the actual disk flush, offering lower latency but risking data loss on unclean shutdowns. The mode is controlled by flush_disk_type in MessageStoreConfig.

How does RocketMQ-Rust ensure data consistency during a broker crash?

The system uses StoreCheckpoint (defined in store_checkpoint.rs) to record the highest committed offsets for the CommitLog and ConsumeQueue. Upon restart, the broker loads this checkpoint and truncates any partial writes beyond the checkpoint offset, ensuring recovery to a known consistent state. Periodic checkpoint flushing integrates with the flush manager pipeline.

What is the role of In-Sync Replicas (ISR) in RocketMQ-Rust durability?

ISR defines how many slave brokers must acknowledge a message before the master considers it committed. Configured via in_sync_replicas in MessageStoreConfig and calculated in CommitLog::calc_need_ack_nums, this quorum-based approach ensures that messages survive even if individual nodes fail, providing distributed durability beyond single-node disk persistence.

How can I configure RocketMQ-Rust for maximum durability?

Set flush_disk_type = FlushDiskType::SyncFlush to ensure every message hits disk before acknowledgment, configure broker_role to SyncMaster to enable HA replication, and set in_sync_replicas to match your replication factor (e.g., 2 for a three-node cluster). Disable enable_auto_in_sync_replicas to prevent the system from lowering the durability requirement during replica lag.

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 →