RocketMQ Cluster Architecture and Advanced Features Guide

1. Cluster Architecture

RocketMQ utilizes a distributed architecture with multiple brokers providing services. The cluster consists of multiple master nodes and slave nodes.

1.1 Cluster Characteristics

NameServer acts as a nearly stateless component that can be deployed in clusters with no information synchronization between nodes.
Broker deployment is relatively complex. Brokers are divided into Master and Slave instances. A Master can correspond to multiple Slaves, but a Slave can only correspond to one Master. The relationship between Master and Slave is defined by specifying the same BrokerName with different BrokerIds. A BrokerId of 0 indicates Master, while non-zero indicates Slave. Multiple Masters can also be deployed. Each Broker establishes long-lived connections with all nodes in the NameServer cluster and periodically registers Topic information.
Producer establishes a long-lived connection with one random node in the NameServer cluster, periodically retrieves Topic routing information from NameServer, and establishes long-lived connections with the Master providing Topic services, periodically sending heartbeats. Producer is completely stateless and can be deployed in clusters.
Consumer establishes a long-lived connection with one random node in the NameServer cluster, periodically retrieves Topic routing information from NameServer, and establishes long-lived connections with the Master and Slave providing Topic services, periodically sending heartbeats. Consumer can subscribe to messages from either Master or Slave, with subscription rules determined by Broker configuration.

Key Points:

  • Multiple NameServer instances available
  • Multiple Masters, each with multiple Slaves
  • Master has brokerId=0, Slave has non-zero brokerId
  • Multiple Masters and Slaves with the same BrokerName form a group
  • Both Master and Slave register with every NameServer instance

1.2 Cluster Workflow

Step 1: NameServer starts and listens for connections from brokers, producers, and consumers
Step 2: Broker starts, connects to all NameServers based on configuration, and maintains long connections
Step 2 (补充): If existing data exists in the broker, NameServer saves the topic-broker relationship
Step 3: Producer sends information, connects to a specific NameServer, and establishes a long connection
Step 4: Producer sends message
    Step 4.1: If topic exists, NameServer directly assigns it
    Step 4.2: If topic doesn't exist, NameServer creates topic-broker relationship and assigns
Step 5: Producer selects a message queue from the list for the topic on the broker
Step 6: Producer establishes long connection with broker for sending messages
Step 7: Producer sends message

Consumer workflow is similar to producer

1.3 Cluster Setup

Step 1: Configure Host Names

vim /etc/hosts

# nameserver
192.168.184.128 rocketmq-nameserver1
192.168.184.129 rocketmq-nameserver2
# broker
192.168.184.128 rocketmq-master1
192.168.184.129 rocketmq-slave2
192.168.184.129 rocketmq-master2
192.168.184.128 rocketmq-slave1

Aply configuration by restarting network:

systemctl restart network

Step 2: Disable Firewall

# Disable firewall
systemctl stop firewalld.service 
# Check firewall status
firewall-cmd --state 
# Prevent firewall from starting on boot
systemctl disable firewalld.service

Step 3: Configure JDK

Refer to Day 01 documentation.

Step 4: Configure Server Environment

Extract RocketMQ to root directory /:

# Extract
unzip rocketmq-all-4.5.2-bin-release.zip
# Rename directory
mv rocketmq-all-4.5.2-bin-release rocketmq

Edit profile:

vim /etc/profile

#set rocketmq
ROCKETMQ_HOME=/rocketmq
PATH=$PATH:$ROCKETMQ_HOME/bin
export ROCKETMQ_HOME PATH

Apply configuration:

source /etc/profile

Step 5: Create Data Storage Directories

Create four directories for master node and four for slave node:

mkdir /rocketmq/store
mkdir /rocketmq/store/commitlog
mkdir /rocketmq/store/consumequeue
mkdir /rocketmq/store/index

mkdir /rocketmq-slave/store
mkdir /rocketmq-slave/store/commitlog
mkdir /rocketmq-slave/store/consumequeue
mkdir /rocketmq-slave/store/index

Note: If Master and Slave are deployed on the same virtual machine, storage directories must be differentiated.

Step 6: Modify Configuration

Different nodes require different configuration files:

cd /rocketmq/conf/2m-2s-sync
vim broker-a.properties  
# Cluster name
brokerClusterName=rocketmq-cluster
# Broker name - different configs should have different values
brokerName=broker-a
# 0 indicates Master, >0 indicates Slave
brokerId=1
# NameServer addresses, semicolon separated
namesrvAddr=rocketmq-nameserver1:9876;rocketmq-nameserver2:9876
# Default queue count for auto-created topics
defaultTopicQueueNums=4
# Allow Broker to auto-create topics (enable offline, disable online)
autoCreateTopicEnable=true
# Allow Broker to auto-create subscription groups (enable offline, disable online)
autoCreateSubscriptionGroup=true
# Broker listening port
listenPort=11011
# File deletion time (default 4 AM)
deleteWhen=04
# File retention time (default 48 hours)
fileReservedTime=48

# CommitLog file size (default 1G)
mapedFileSizeCommitLog=1073741824
# ConsumeQueue file stores 300K messages by default
mapedFileSizeConsumeQueue=300000
# Maximum disk space usage ratio
diskMaxUsedSpaceRatio=88
# Storage path
storePathRootDir=/rocketmq/store-slave
# CommitLog storage path
storePathCommitLog=/rocketmq/store-slave/commitlog
# ConsumeQueue storage path
storePathConsumeQueue=/rocketmq/store-slave/consumequeue
# Message index storage path
storePathIndex=/rocketmq/store-slave/index
# Checkpoint file storage path
storeCheckpoint=/rocketmq/store-slave/checkpoint
# Abort file storage path
abortFile=/rocketmq

# Maximum message size
maxMessageSize=65536
# Broker role
# - ASYNC_MASTER: Async replication Master
# - SYNC_MASTER: Sync dual-write Master
# - SLAVE
brokerRole=SLAVE
# Flush type
# - ASYNC_FLUSH: Async flush
# - SYNC_FLUSH: Sync flush
flushDiskType=SYNC_FLUSH
# Message send thread pool size
#sendMessageThreadPoolNums=128
# Message pull thread pool size
#pullMessageThreadPoolNums=128

Adjust startup memory (modify both nameserver and broker):

vim /rocketmq/bin/runbroker.sh
vim /rocketmq/bin/runserver.sh

# Development environment JVM configuration
JAVA_OPT="${JAVA_OPT} -server -Xms256m -Xmx256m -Xmn128m"

Start services (from bin directory):

nohup sh mqnamesrv &
nohup sh mqbroker -c ../conf/2m-2s-sync/broker-a.properties &
nohup sh mqbroker -c ../conf/2m-2s-sync/broker-b-s.properties &

2. RocketMQ Console

RocketMQ Console is a Java-based (Spring Boot) management console tool. Obtain from: https://github.com/apache/rocketmq-externals

3. Advanced Features

3.1 Message Storage

ActiveMQ uses database storage, which creates a bottleneck at the database level. RocketMQ, Kafka, and RabbitMQ directly use file storage instead.

3.2 Efficient Message Storage and Read/Write

1) Pre-initialize file sizes at startup to reserve fixed disk space and ensure disk I/O performance
2) Zero-copy technology: Data transfer is reduced from 4 copies to 3 copies, eliminating one copy operation. Java's MappedByteBuffer implements this technology. Requirement: Reserve storage space for data (starting from 1G)

3.3 Message Storage Structure

Message Data Storage Area
    - topic
    - queueId
    - message

Consume Logic Queue
    - minOffset
    - maxOffset
    - consumerOffset

Index
    - key index
    - creation time index
    - ...

3.4 Flush Mechanism

Synchronous Flush

1) Producer sends message to MQ, MQ receives message data
2) MQ suspends the producer's sending thread
3) MQ writes message data to memory
4) Memory data is written to disk
5) After disk storage, return SUCCESS
6) MQ resumes the suspended producer thread
7) Send ACK to producer

Asynchronous Flush

1) Producer sends message to MQ, MQ receives message data
2) MQ writes message data to memory
3) Send ACK to producer
-- When message volume accumulates --
4) Memory data is written to disk

Comparison: Synchronous vs Asynchronous Flush

Synchronous Flush: High safety, low efficiency, slow speed (suitable for businesses requiring high data security)
Asynchronous Flush: Low safety, high efficiency, fast speed (suitable for businesses requiring high processing speed)

Configuration

# Flush type
# - ASYNC_FLUSH: Async flush
# - SYNC_FLUSH: Sync flush
flushDiskType=SYNC_FLUSH

4. High Availability

NameServer

NameServer ensures service continuity through stateless architecture with full-server registration, allowing service provision even if one instance fails.

Message Server

Master-Slave architecture (2M-2S) ensures service availability even when one server fails. Note: When Master fails, Slave only provides consumption services and cannnot accept new messages (Slave does not upgrade to Master).

Message Production

Developers should bind the same topic to multiple group IDs to ensure other Masters can receive messages normally when one Master fails.

Message Consumption

RocketMQ automatically determines whether Master or Slave handles message reading based on Master load. When Master is busy, it automatically switches to Slave for data reading.

5. Master-Slave Data Replication

Synchronous Replication

Master replicates message to Slave first, then responds to producer with write success.

  • Advantages: Data is safe, no data loss, easy recovery from failures
  • Disadvantages: Affects data throughput, lower overall performance

Asynchronous Replication

Master immediately returns write success to producer, then asynchronously replicates to Slave when message volume accumulates.

  • Advantages: High data throughput, low latency, high performance
  • Disadvantages: Data is not safe, data loss may occur; data from last sync to failure time will be lost

Configuration (in broker config file specified by -c parameter)

# Broker role
# - ASYNC_MASTER: Async replication Master
# - SYNC_MASTER: Sync dual-write Master
# - SLAVE
brokerRole=SYNC_MASTER

6. Load Balancing

Producer Load Balance

Internally implements load balancing across message queues of the same topic across different broker clusters.

Consumer Load Balance

  • Average allocation
  • Round-robin average allocation

7. Message Retry

Message retry mechanism activates when consumer fails to return successful consumption status.

Ordered Message Retry

When consumer fails to process a message, RocketMQ automatically retries (interval: 1 second)
Note: Application may experience message consumption blocking; monitor ordered message consumption to prevent blocking

Unordered Message Retry

Unordered messages include regular messages, scheduled messages, delayed messages, and transaction messages
Unordered message retry only applies to load-balanced (cluster) consumption model, not broadcast mode
MQ sets reasonable retry intervals for unordered message consumption

8. Dead Letter Queue

Concept

When message retry reaches the specified count (default 16 times), messages that cannot be consumed normally are called Dead-Letter Messages. These are not discarded but saved to a new queue called Dead-Letter Queue.

Dead Letter Queue Characteristics

  • Belongs to a specific Group ID, not Topic or Consumer
  • Can contain dead letter messages from multiple Topics under the same group
  • Not initialized by default; first dead letter triggers initialization

Dead Letter Message Characteristics

  • Will not be consumed again
  • Messages remain for 3 days, then deleted

Handling Dead Letters

In monitoring platform, find dead letters by messageId and precisely consume them.

9. Message Duplicate Consumption and Idempotency

Causes of Message Duplicate Consumption

1) Producer sends duplicate messages
    - Network interruption
    - Producer crash
2) Message server delivers duplicate messages
    - Network interruption
3) Dynamic load balancing process
    - Network interruption/jitter
    - Broker restart
    - Consumer application restart
    - Client scale-up
    - Client scale-down

Message Idempotency

For the same message, regardless of how many times it is consumed, the result remains consistent—this is message idempotency.

Solution

  • Use business ID as message key
  • Consumer checks the key; if unused, allow processing; if used, discard

Note: messageId is generated by RocketMQ and is not unique, therefore cannot be used as idempotency判定 condition.

Tags: rocketmq Message Queue Cluster Architecture Distributed Systems MQ

Posted on Thu, 01 Oct 2026 16:38:55 +0000 by DwarV