# How RocketMQ-Rust Handles Failure Recovery: Architecture and Implementation

> RocketMQ-Rust ensures automatic failure recovery through robust architecture and implementation strategies. Learn how optimized iterators and memory mapped file validation restore broker state.

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

---

**RocketMQ-Rust implements automatic failure recovery across the CommitLog, ConsumeQueue, and high-availability layers, using optimized batch iterators, memory-mapped file validation, and configurable retry logic to restore broker state after crashes or network partitions.**

RocketMQ-Rust (mxsm/rocketmq-rust) is a Rust implementation of the Apache RocketMQ broker that maintains data consistency through comprehensive recovery mechanisms. This article examines how the codebase handles failure recovery across storage, indexing, and high-availability components using idiomatic Rust patterns and zero-copy optimizations.

## CommitLog Recovery Mechanisms

The CommitLog serves as the primary storage layer for message data. RocketMQ-Rust provides two distinct recovery paths 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) and [`rocketmq-store/src/log_file/commit_log_recovery.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/log_file/commit_log_recovery.rs): optimized batch recovery and normal sequential recovery.

### Optimized Batch Recovery

When the environment variable `ROCKETMQ_USE_OPTIMIZED_RECOVERY` is set to `"true"` (the default), the broker invokes `recover_normally_optimized` or `recover_abnormally_optimized`. These methods utilize a **zero-copy batch iterator** (`BatchMessageIterator`) to read commit-log files in **64 KB chunks**, rebuilding the physical offset map and truncating dirty data. The iterator validates each message via `check_message_and_return_size` while updating `last_valid_msg_phy_offset`.

### Normal Sequential Recovery

If optimized recovery is disabled, the system falls back to `recover_normally` or `recover_abnormally`, which walk the log sequentially without batching. Both paths ultimately dispatch valid messages to the consumer queue through the `Dispatcher` interface when replay is required.

## ConsumeQueue and Index Recovery

After CommitLog recovery completes, each `ConsumeQueue` instance (implemented in [`rocketmq-store/src/queue/single_consume_queue.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/queue/single_consume_queue.rs) and related batch queue files) executes its `recover` method. This process iterates over memory-mapped files, validates offsets against the restored CommitLog, and truncates corrupted tails to ensure index consistency.

## Topic-Queue Table Reconstruction

The broker restores its in-memory routing metadata through `recover_topic_queue_table` 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). This function walks the `ConsumeQueue` metadata after CommitLog recovery and rebuilds the hash table mapping `topic → queueId → offset`, enabling proper message routing immediately after startup.

## High Availability Client Recovery

For Slave brokers, the `DefaultHAClient` (located in [`rocketmq-store/src/ha/default_ha_client.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/default_ha_client.rs)) implements automatic reconnection logic. When the master-slave HA link breaks due to network partitions or master restarts, the client enters a retry loop that attempts to open a TCP socket to the master every **5 seconds** (`sleep(Duration::from_secs(5)).await`). Upon successful connection, it reinitializes replication state via `change_current_state` and transmits a `HAConnectionStateNotificationRequest`.

## HA Service State Machine and Offset Synchronization

The `DefaultHAService` (in [`rocketmq-store/src/ha/default_ha_service.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/default_ha_service.rs)) manages the state machine transitions (Ready → Transfer → Slave) and handles concurrent offset updates. The service stores the latest reported offset in an `AtomicU64` and applies updates using a **compare-and-swap loop** (`compare_exchange_weak`) that retries on contention, ensuring thread-safe progress tracking during replication recovery.

## Scheduled Task Retry Logic

For operational tasks that may fail transiently, [`rocketmq/src/schedule/task.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq/src/schedule/task.rs) implements bounded retry logic. Each `Task` carries a `max_retry` field; the scheduler increments `retry_count` on each failure and terminates the task when the limit is reached. This prevents infinite loops while allowing temporary faults to resolve automatically.

## Broker Startup Recovery Flow

The `LocalFileMessageStore` orchestrates a six-phase recovery sequence during broker initialization:

1. **Component Initialization**: `LocalFileMessageStore::new` creates the HA client (for Slave roles) and instantiates `CommitLog` and `ConsumeQueue` structures.

2. **Store Recovery Invocation**: `LocalFileMessageStore::recover` delegates to `recover_consume_queue`, which selects between `recover_normally` and `recover_abnormally` based on shutdown status.

3. **CommitLog Processing**: The system chooses the optimized path when `ROCKETMQ_USE_OPTIMIZED_RECOVERY=true`. The batch iterator reads each mapped file, validates messages, updates `last_valid_msg_phy_offset`, and dispatches to the `Dispatcher`.

4. **ConsumeQueue Replay**: If `broker_config.recover_concurrently` is enabled, ConsumeQueue files are replayed concurrently. Each queue truncates dirty tails based on the final CommitLog offset from phase 3.

5. **Metadata Rebuild**: The **Topic-Queue table** is reconstructed from the validated ConsumeQueue data.

6. **HA Service Activation**: For Slave brokers, the background reconnect loop starts automatically, repeatedly calling `connect_master` with 5-second back-off until the master responds.

## Configuration and Usage Examples

### Enabling Optimized Recovery

Set the environment variable before starting the broker to use the 64 KB batch iterator:

```rust
std::env::set_var("ROCKETMQ_USE_OPTIMIZED_RECOVERY", "true");

```

The broker automatically selects the optimized path in `LocalFileMessageStore::recover_normally` or `recover_abnormally`.

### Configuring HA Client Reconnection

When initializing a Slave broker, the HA client runs autonomously:

```rust
use rocketmq_store::ha::default_ha_client::DefaultHAClient;
use rocketmq_store::message_store::local_file_message_store::LocalFileMessageStore;
use std::sync::Arc;

let ha_client = DefaultHAClient::new(msg_store.clone())
    .expect("Failed to create HA client");

tokio::spawn(async move {
    ha_client.start().await;
});

```

### Defining Retry Limits for Scheduled Tasks

Configure task resilience using the `max_retry` parameter:

```rust
use rocketmq::schedule::{Task, TaskContext, TaskResult};
use std::time::Duration;

let my_task = Task::new("id-1", "cleanup", |ctx: TaskContext| async move {
    // Task implementation
    TaskResult::Ok
})
.with_max_retry(3)
.with_initial_delay(Duration::from_secs(10))
.enabled(true);

```

## Summary

- **RocketMQ-Rust** provides dual-path CommitLog recovery: an optimized zero-copy batch mode using `BatchMessageIterator` (64 KB chunks) and a fallback sequential mode.
- **ConsumeQueue** files recover independently by validating offsets against the restored CommitLog and truncating corrupted tails.
- The **Topic-Queue table** rebuilds automatically from clean ConsumeQueue metadata via `recover_topic_queue_table`.
- **HA clients** implement exponential-less retry logic with fixed 5-second intervals to reconnect Slaves to Masters after network partitions.
- **State synchronization** uses atomic compare-and-swap operations (`compare_exchange_weak`) to manage replication offsets safely across threads.
- **Scheduled tasks** support bounded retries via the `max_retry` field in `Task` definitions.

## Frequently Asked Questions

### How does RocketMQ-Rust detect that recovery is needed?

The broker triggers recovery automatically during `LocalFileMessageStore::new` initialization or when detecting an unclean shutdown. The `recover` method checks the commit-log state and invokes `recover_abnormally` if the previous shutdown was not graceful, otherwise it calls `recover_normally`.

### What is the performance difference between optimized and normal recovery?

The optimized path in [`commit_log_recovery.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/commit_log_recovery.rs) uses a `BatchMessageIterator` to process messages in 64 KB chunks with zero-copy memory mapping, significantly reducing I/O overhead compared to the normal path, which walks the log sequentially without batching.

### How does the HA client handle persistent connection failures?

The `DefaultHAClient` (in [`rocketmq-store/src/ha/default_ha_client.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-store/src/ha/default_ha_client.rs)) implements an infinite retry loop with a fixed 5-second back-off (`Duration::from_secs(5)`). It attempts to `connect_master` repeatedly, sleeping between attempts, until the TCP connection succeeds and replication resumes.

### Can ConsumeQueue recovery run in parallel?

Yes. When `broker_config.recover_concurrently` is enabled, the broker replays multiple ConsumeQueue files concurrently during startup. Each queue independently validates its offsets against the final CommitLog position and truncates dirty data accordingly.