How RocketMQ-Rust Achieves High Throughput and Low Latency: 8 Core Optimizations Explained
RocketMQ-Rust achieves high throughput and low latency through zero-copy encoding buffers, lock-free message processing with narrow critical sections, object-pool reuse for encoders, and a customized Tokio runtime that enables asynchronous I/O with per-queue sharding.
RocketMQ-Rust (mxsm/rocketmq-rust) is a Rust rewrite of Apache RocketMQ designed to handle massive message streams with minimal latency. The implementation combines lock-free algorithms, adaptive memory management, and fine-grained concurrency controls to deliver throughput exceeding 10,000 TPS with P99 latencies under 25 milliseconds, as verified by the commit_log_performance_tests.rs integration suite.
Lock-Free Encoding and Minimal Critical Sections
The broker eliminates contention by performing message encoding before acquiring any locks. In rocketmq-store/src/message_store/local_file_message_store.rs, the encoding phase runs lock-free, ensuring that the critical section for the per-topic-queue lock covers only offset assignment (approximately 0.1ms).
This narrow lock scope means that when multiple producers write to different queues concurrently, they do not block each other. The lock protects only the offset-assignment step; serialization, buffer handling, and I/O operations execute outside the critical section. This architecture enables eight or more queues to be written in parallel without cross-queue blocking.
Zero-Copy Adaptive Buffering with EncodeBuffer
RocketMQ-Rust implements an adaptive buffering strategy in rocketmq-remoting/src/smart_encode_buffer.rs to eliminate memory waste and reallocation overhead. The EncodeBuffer type wraps BytesMut and grows on demand while shrinking conservatively based on an exponential moving average (EMA) of recent write sizes.
After a cooldown period, oversized buffers automatically release excess capacity, preventing long-lived memory bloat while maintaining hot-path performance. This zero-copy approach allows encoded payloads to move directly to the network layer without intermediate allocations.
use rocketmq_remoting::smart_encode_buffer::EncodeBuffer;
use bytes::BytesMut;
/// Encode a payload without allocating a new buffer each time.
fn encode_payload(payload: &[u8]) -> BytesMut {
// Buffer starts with a modest 8 KB capacity.
let mut eb = EncodeBuffer::new();
eb.append(payload); // Expands only if needed.
let bytes = eb.take_bytes(); // Zero‑copy Bytes.
// `bytes` can be handed to the network layer directly.
bytes.into()
}
Memory Optimization via Object Pooling
To cut allocation overhead by approximately 50%, the store uses encoder key pooling via generate_key_with_pool in rocketmq-store/src/base/message_encoder_pool.rs. Instead of creating new encoder instances for every message, the system reuses keys across messages sharing the same topic and queue, significantly reducing heap pressure during high-throughput scenarios.
use rocketmq_store::base::message_encoder_pool::generate_key_with_pool;
use rocketmq_common::common::message::message_ext_broker_inner::MessageExtBrokerInner;
// Assume `msg` is a MessageExtBrokerInner.
let encoder_key = generate_key_with_pool(&msg);
// The same `encoder_key` will be reused for subsequent messages
// sharing the same topic/queue, avoiding new allocations.
Asynchronous Runtime and I/O Tuning
All network and disk operations run on a customized Tokio runtime defined in rocketmq-common/src/thread_pool.rs. The TokioExecutorService provides configurable worker threads, keep-alive durations, and maximum blocking threads, allowing the broker to scale with CPU core count while maintaining low latency.
Background maintenance tasks—such as cleaning and statistics collection—execute on a separate ScheduledExecutorService with fixed-rate guarantees, ensuring that periodic work never interferes with the critical request path.
use rocketmq_common::thread_pool::TokioExecutorService;
use std::time::Duration;
let executor = TokioExecutorService::new_with_config(
8, // 8 worker threads
Some("mq-worker-".to_string()),
Duration::from_secs(60), // keep‑alive
256, // max blocking threads
);
executor.spawn(async move {
// Periodic stats collection
loop {
// …collect and log stats…
tokio::time::sleep(Duration::from_secs(5)).await;
}
});
Batch Processing and Storage Optimization
The local_file_message_store.rs module writes messages in batches using MessageExtBatch and pre-allocates commit-log files to eliminate latency spikes caused by on-the-fly file growth. The store detects full batches efficiently and flushes them using optimized logic.
The commit-log flush implementation uses a single match expression over the four combinations of sync/async flush and HA/no-HA configurations. This removes duplicated branching and ensures the fastest path is taken for each durability requirement, minimizing code-path length during the hot flush cycle.
Concurrent Multi-Queue Architecture
Because each queue maintains its own independent lock, the broker scales horizontally across queues. This per-queue isolation allows the system to sustain high throughput even under heavy concurrent load, with performance tests in rocketmq-store/tests/commit_log_performance_tests.rs demonstrating sustained throughput above 10,000 TPS across multiple concurrent queues.
Summary
- Lock-free encoding moves serialization outside critical sections, reducing the Topic-Queue lock duration to approximately 0.1ms.
- Adaptive EncodeBuffer with EMA-driven shrinking prevents memory bloat while maintaining zero-copy network transfers.
- Object pooling via
generate_key_with_poolreduces encoder allocation overhead by roughly 50%. - Custom Tokio runtime (
TokioExecutorService) scales with CPU cores and isolates blocking operations. - Batch processing with pre-allocated commit-log files eliminates file-growth latency spikes.
- Per-queue locking enables parallel writes across eight or more queues without cross-queue contention.
Frequently Asked Questions
How does RocketMQ-Rust minimize lock contention during message writes?
RocketMQ-Rust performs message encoding before acquiring the Topic-Queue lock, keeping the critical section limited to offset assignment only. This narrow lock scope, combined with per-queue locking in the store layer, allows concurrent writes to different queues without blocking each other.
What buffer management strategy prevents memory bloat in high-throughput scenarios?
The EncodeBuffer in rocketmq-remoting/src/smart_encode_buffer.rs uses BytesMut with an exponential moving average (EMA) of recent write sizes to shrink buffers conservatively after a cooldown period. This approach prevents long-lived memory bloat while avoiding frequent reallocations during encoding.
How does the Tokio runtime configuration impact latency and throughput?
The TokioExecutorService configures worker threads proportional to CPU cores, sets explicit limits on blocking threads (defaulting to 256), and manages keep-alive durations. This tuning ensures that asynchronous network and disk I/O remain non-blocking while dedicated threads handle periodic background tasks via ScheduledExecutorService without interfering with the hot path.
How does RocketMQ-Rust handle message encoding without excessive allocation overhead?
The system uses generate_key_with_pool from rocketmq-store/src/base/message_encoder_pool.rs to reuse encoder keys across messages with the same topic and queue. This object-pooling strategy cuts allocation overhead by approximately 50%, as measured by the store's memory-allocation regression tests.
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 →