RocketMQ Rust Client Features: Producer and Consumer Capabilities Explained

The rocketmq-client crate provides a full-featured, async-first implementation of RocketMQ messaging primitives for Rust, exposing DefaultMQProducer for granular message publishing with configurable retry, compression, and back-pressure policies, and DefaultMQPushConsumer for subscription-based consumption with automatic rebalancing and flow control.

The mxsm/rocketmq-rust repository delivers a native Rust client for Apache RocketMQ that prioritizes asynchronous operations throughout its public API. Whether building high-throughput event streaming pipelines or request-reply microservices, developers interact with two primary high-level abstractions: DefaultMQProducer for message publishing and DefaultMQPushConsumer for message consumption. Both components expose extensive configuration options via builder patterns defined in rocketmq-client/src/producer/default_mq_produce_builder.rs and rocketmq-client/src/consumer/default_mq_push_consumer_builder.rs.

Producer Features and Configuration

Initialization and Lifecycle

Producers are instantiated using the ergonomic builder API exposed by DefaultMQProducer::builder(). According to the source in rocketmq-client/src/producer/default_mq_producer.rs lines 78-80, the builder accepts critical parameters including producer_group, send_msg_timeout, and retry_times_when_send_failed.

The start() method initializes the internal DefaultMQProducerImpl and optionally registers distributed tracing hooks when client_config.enable_trace is true (lines 94-104). Clean resource management is ensured via the shutdown() method, which terminates background threads and closes network connections.

use rocketmq_client::producer::{DefaultMQProducer, MQProducer};

let mut producer = DefaultMQProducer::builder()
    .producer_group("order-service")
    .send_msg_timeout(5000)
    .retry_times_when_send_failed(3)
    .auto_batch(true)
    .build();

producer.start().await?;

Message Sending Patterns

The producer supports three distinct delivery semantics. Synchronous sending via send(msg) blocks until the broker acknowledges receipt (lines 55-64). Asynchronous sending via send_with_callback(msg, callback) returns immediately and executes the provided closure upon completion or failure (lines 84-99). For fire-and-forget scenarios, send_oneway(msg) transmits messages without waiting for acknowledgment.

Targeting specific partitions is supported through send_to_queue(msg, mq), allowing deterministic routing to a MessageQueue instance.

Advanced Capabilities: Batching, Selectors, and Compression

High-throughput scenarios benefit from batch sending via send_batch(msgs), which constructs a MessageBatch internally and transmits multiple messages in a single network round-trip (lines 52-64). The batch automatically respects broker size limits and handles partial failures.

For custom routing logic, send_with_selector(msg, selector, arg) accepts a user-defined closure implementing the MessageQueueSelector trait (lines 126-135). This enables affinity-based routing or sharding strategies.

Compression is handled transparently using ZLIB when message bodies exceed the compress_msg_body_over_howmuch threshold configured in ProducerConfig. Developers may inject custom compression algorithms via the compressor extensibility hook.

Back-Pressure and Fault Tolerance

The producer implements sophisticated back-pressure mechanisms for async operations. When enable_backpressure_for_async_mode is true, the client monitors pending requests via back_pressure_for_async_send_num and back_pressure_for_async_send_size, throttling new sends when thresholds are exceeded (lines 100-110).

Retry policies are configurable through retry_response_codes and retry_another_broker_when_not_store_ok, allowing automatic failover to secondary brokers when the primary rejects messages. Tracing integration via AsyncTraceDispatcher registers SendMessageTraceHookImpl hooks to emit distributed tracing spans without blocking the send path (lines 106-122).

Consumer Features and Architecture

Push Consumer Configuration

The DefaultMQPushConsumer manages subscription-based message consumption through the builder pattern defined in rocketmq-client/src/consumer/default_mq_push_consumer.rs lines 31-42. Configuration options include consumer_group, message_model (clustering or broadcasting), and allocate_message_queue_strategy for partition assignment.

The start() method initializes DefaultMQPushConsumerImpl, wraps the consumer group with the current namespace, and registers ConsumeMessageTraceHookImpl when tracing is enabled (lines 22-30).

use rocketmq_client::consumer::{DefaultMQPushConsumer, MQConsumer};

let mut consumer = DefaultMQPushConsumer::builder()
    .consumer_group("inventory-service")
    .subscribe("order-events", "*")
    .build();

consumer.start().await?;

Subscription Management

Consumers declare interest in topics via subscribe(topic, expression), where the expression specifies tag-based filtering (lines 88-100). Dynamic subscription changes are supported through unsubscribe(topic). During maintenance windows, suspend() pauses message fetching without shutting down the consumer, while resume() restores normal operation (lines 17-27).

Message Processing and Flow Control

Message consumption relies on registered listeners implementing MessageListenerConcurrently or MessageListenerOrderly. The listener receives a Vec<MessageExt>, allowing batch processing of messages fetched from the broker. The consumer does not expose a standalone "batch receive" API; instead, batching occurs internally, and the user-provided callback handles message collections.

Flow control is configurable through ConsumerConfig parameters including pull_threshold_for_queue, consume_thread_count, and max_reconsume_times. When thresholds are exceeded, the consumer automatically reduces pull frequency to prevent memory exhaustion. The pluggable AllocateMessageQueueStrategy determines partition assignment during consumer group rebalances.

Practical Implementation Examples

Implementing a Push Consumer with Concurrent Processing

The following example demonstrates subscribing to a topic and processing messages concurrently using async closures:

use rocketmq_client::consumer::{
    DefaultMQPushConsumer, MQConsumer, MessageListenerConcurrently,
    listener::ConsumeConcurrentlyStatus,
};
use rocketmq_common::common::message::message_ext::MessageExt;

let mut consumer = DefaultMQPushConsumer::builder()
    .consumer_group("log-processor")
    .subscribe("application-logs", "ERROR || WARN")
    .build();

consumer.register_message_listener_concurrently(
    |msgs: Vec<MessageExt>| async move {
        for msg in msgs {
            let payload = std::str::from_utf8(msg.get_body()).unwrap_or("invalid utf8");
            println!("[{}] {}: {}", msg.get_msg_id(), msg.get_topic(), payload);
        }
        ConsumeConcurrentlyStatus::ConsumeSuccess
    },
);

consumer.start().await?;
tokio::signal::ctrl_c().await?;
consumer.shutdown().await;

Custom Queue Selection for Producers

When message ordering or partition affinity is required, use send_with_selector to implement custom routing logic:

let selector = |queues: &[MessageQueue], _msg: &Message, arg: &i32| {
    queues.get((*arg as usize) % queues.len()).cloned()
};

producer
    .send_with_selector(
        Message::builder()
            .topic("user-events")
            .body(Bytes::from_static(b"user-123-action"))
            .build_unchecked(),
        selector,
        123, // user_id used for routing
    )
    .await?;

Request-Reply Pattern

The client supports synchronous request-reply semantics via the request method, which waits for a consumer to process the message and return a correlated response:

let reply = producer
    .request(
        Message::builder()
            .topic("rpc-queue")
            .body(Bytes::from_static(b"ping"))
            .build_unchecked(),
        3000, // timeout in milliseconds
    )
    .await?;

println!("Received reply: {:?}", reply);

Summary

  • The rocketmq-client crate exposes async-first APIs through DefaultMQProducer and DefaultMQPushConsumer, with comprehensive configuration via builder patterns in default_mq_produce_builder.rs and default_mq_push_consumer_builder.rs.
  • Producers support multiple send semantics (sync, async, oneway), batch transmission via send_batch, and deterministic routing through send_with_selector.
  • Built-in back-pressure mechanisms and configurable retry policies with broker failover ensure resilient message publishing under high load.
  • Push consumers manage subscriptions through subscribe(), process messages via concurrent or orderly listeners, and automatically handle partition rebalancing using pluggable allocation strategies.
  • Distributed tracing is integrated via AsyncTraceDispatcher for both producers and consumers, while compression and RPC hooks provide additional extensibility points.

Frequently Asked Questions

What is the difference between synchronous and asynchronous message sending?

Synchronous sending via send() blocks the async task until the broker acknowledges receipt, returning a SendResult containing the message ID and offset. Asynchronous sending via send_with_callback() returns immediately and executes a user-provided closure upon completion, enabling higher throughput by decoupling message transmission from result handling. The callback receives either the SendResult or an error explaining the failure reason.

How does the push consumer handle message redistribution when consumers join or leave the group?

The DefaultMQPushConsumer uses the AllocateMessageQueueStrategy configured in ConsumerConfig to determine partition ownership during rebalancing. When group membership changes, the internal DefaultMQPushConsumerImpl triggers a rebalance algorithm that revokes and assigns message queues among active instances. This process updates local consumption offsets and ensures exactly-once processing semantics per partition within the consumer group.

Is transaction support available in the Rust rocketmq-client?

Transaction support is currently stubbed but not implemented. The send_message_in_transaction method in default_mq_producer.rs (lines 40-50) contains an unimplemented!() placeholder, indicating that half-messages and commit/rollback coordination are not yet available. For atomic operations, consider implementing compensating transactions at the application layer or using the request-reply pattern for synchronous confirmation.

How can I implement flow control to prevent consumer overload?

Configure flow control through ConsumerConfig parameters before building the consumer. Set pull_threshold_for_queue to limit the number of cached messages per queue, and adjust consume_thread_count to constrain concurrent processing. When thresholds are exceeded, the consumer automatically reduces pull frequency, creating back-pressure that propagates to the broker. Additionally, implementing rate limiting within your MessageListenerConcurrently callback provides application-level protection against traffic spikes.

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 →