rocketmq-controller Status and Functionality in the RocketMQ Rust Ecosystem
The rocketmq-controller is a high-availability controller component currently in active development that implements Raft consensus for managing cluster metadata and broker lifecycle in the RocketMQ Rust ecosystem.
The rocketmq-controller crate within the mxsm/rocketmq-rust repository provides the core controller functionality for the Rust implementation of Apache RocketMQ. This component serves as the central coordination service for cluster metadata, ensuring strong consistency across distributed broker nodes through the Raft consensus algorithm.
Current Development Status of rocketmq-controller
Repository Status
According to the monorepo table in the top-level README, rocketmq-controller is marked with 🚧 In Development. While the API is still stabilizing, the core functionality—including Raft-based leader election, log replication, and broker metadata management—is already implemented and usable for internal testing and early-stage deployments.
Stability and API State
The crate follows semantic versioning within the broader mxsm/rocketmq-rust workspace. Key structures like ControllerConfig and ControllerCli are stabilized in src/config.rs and src/cli.rs respectively, though the internal Raft state machine interfaces may undergo refinement as the integration with the open-raft library matures.
Core Functionality of rocketmq-controller
Cluster Metadata Management
The controller maintains the single source of truth for all cluster-wide metadata. This includes broker registrations, topic configurations, consumer group offsets, and cluster topology. When a broker starts, it registers with the controller via the register_broker processor, which stores the broker's identity and endpoint information in the Raft-backed state machine defined in src/metadata/broker.rs.
Leader Election and Failover
Raft consensus powers the controller's high-availability guarantees. The controller runs as a cluster of nodes (typically 3 or 5) where one node is elected as the leader. If the leader fails, the Raft algorithm automatically triggers a new election. This implementation in src/raft/raft_controller.rs wraps the open-raft library to provide RaftNetwork and RaftStorage implementations tailored to RocketMQ's metadata model.
Raft Log Replication
All state changes—such as broker registration updates or topic configuration modifications—are written to the Raft log before being applied to the state machine. The log replication mechanism ensures that committed entries are durably stored on a majority of controller nodes before acknowledging the operation to clients. This provides the strong consistency required for critical metadata operations.
Broker Heartbeat Management
The controller tracks broker liveness through a heartbeat mechanism implemented in src/heartbeat/default_broker_heartbeat_manager.rs. Brokers periodically send heartbeat requests to the controller leader. If a broker misses a configurable number of heartbeat intervals, the controller marks it as offline and triggers metadata cleanup procedures to remove stale broker entries from the cluster view.
Architecture and Implementation Details
Controller Manager
The ControllerManager struct in src/controller/controller_manager.rs serves as the orchestration layer. It initializes the Raft controller, starts the RPC server for broker communication, and manages background services including the heartbeat monitor and broker housekeeping tasks. The manager implements the start_blocking() method that runs the controller's main event loop.
Raft Controller Implementation
The RaftController in src/raft/raft_controller.rs bridges the open-raft library with RocketMQ-specific storage and networking. It implements:
- RaftNetwork: Handles RPC communication between controller nodes for Raft message passing (AppendEntries, RequestVote).
- RaftStorage: Provides persistent storage for the Raft log and state machine snapshots, typically backed by RocksDB or file-based storage.
Metadata Storage
Metadata structures are defined in src/metadata/ modules:
broker.rs: Broker identity, endpoints, and registration info.topic.rs: Topic configuration and replica assignments.controller_config.rs: Runtime configuration parameters.
These structures implement serialization traits for Raft log entry persistence and gRPC/JSON communication with brokers.
Building and Running rocketmq-controller
To build the controller binary from the mxsm/rocketmq-rust repository:
# Build the release binary
cargo build --release --bin rocketmq-controller-rust
# Start with default configuration (expects $ROCKETMQ_HOME/conf/controller.toml)
./target/release/rocketmq-controller-rust
# Start with explicit config file
./target/release/rocketmq-controller-rust -c controller.toml
# Print effective configuration without starting
./target/release/rocketmq-controller-rust -c controller.toml -p
The command-line interface is defined in src/cli.rs, which validates arguments and loads the ControllerConfig structure. The -p flag triggers the print_config method to display the full effective configuration.
Programmatic Usage
You can embed the rocketmq-controller functionality in your own Rust applications:
use rocketmq_controller::cli::ControllerCli;
use rocketmq_controller::config::ControllerConfig;
use rocketmq_error::RocketMQResult;
fn main() -> RocketMQResult<()> {
// Parse command-line arguments
let cli = ControllerCli::parse_args();
// Validate and load configuration
cli.validate()?;
let default_cfg = ControllerConfig::default();
let cfg = cli.load_config(default_cfg)?;
// Handle configuration printing
if cli.print_config_item {
ControllerCli::print_config(&cfg);
return Ok(());
}
// Start the controller manager
let manager = rocketmq_controller::controller::ControllerManager::new(cfg);
manager.start_blocking()?;
Ok(())
}
This pattern mirrors the actual bootstrap sequence in src/bin/controller_bootstrap.rs. For advanced use cases, you can directly instantiate the OpenRaftController from src/openraft.rs to customize the Raft network or storage layers.
Summary
- rocketmq-controller is the high-availability metadata controller for the RocketMQ Rust implementation, currently in development but functional for testing.
- It implements Raft consensus (
open-raftlibrary) for leader election, log replication, and strong consistency across controller nodes. - Core responsibilities include cluster metadata management, broker registration, heartbeat monitoring, and configuration persistence.
- Key source files include
src/controller/controller_manager.rsfor orchestration,src/raft/raft_controller.rsfor consensus logic, andsrc/heartbeat/default_broker_heartbeat_manager.rsfor broker liveness tracking. - The controller can be run via the
rocketmq-controller-rustbinary with CLI options defined insrc/cli.rs, or embedded programmatically via theControllerManagerAPI.
Frequently Asked Questions
Is rocketmq-controller production-ready?
No, rocketmq-controller is currently marked as 🚧 In Development in the mxsm/rocketmq-rust repository. While the core Raft consensus implementation and metadata management functionality are complete and testable, the API is still stabilizing. It is suitable for internal testing and early-stage deployments but not yet recommended for production workloads without thorough validation.
How does rocketmq-controller differ from RocketMQ's original Java controller?
The rocketmq-controller is a ground-up Rust implementation of the controller logic found in Apache RocketMQ's Java codebase. While both use Raft for consensus, the Rust version leverages the open-raft library and integrates with the broader rocketmq-rust ecosystem. The architecture in src/controller/controller_manager.rs and src/raft/raft_controller.rs mirrors the Java controller's responsibilities but with Rust-specific memory safety and concurrency patterns.
What consensus algorithm does rocketmq-controller use?
rocketmq-controller implements the Raft consensus algorithm through the open-raft library. This provides distributed consensus for leader election, log replication, and cluster membership changes. The implementation in src/raft/raft_controller.rs wraps open-raft and provides the RaftNetwork and RaftStorage traits necessary for RocketMQ-specific metadata persistence and inter-node communication.
How do I configure rocketmq-controller for a multi-node cluster?
Configuration is handled through the ControllerConfig structure, typically loaded from a TOML, JSON, or YAML file specified via the -c flag. For a multi-node cluster, you must define the raft_peers configuration with the node IDs and addresses of all controller instances. Each node should have a unique ID and shared knowledge of the peer set. The ControllerCli in src/cli.rs validates this configuration before starting the ControllerManager, ensuring all Raft nodes can communicate for quorum establishment.
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 →