本文系统梳理消息队列(MQ)的运维知识体系,覆盖 RabbitMQ、Kafka、RocketMQ 三大主流中间件的安装部署、集群管理、消息可靠性保障、性能优化与故障排查,适用于生产环境的日常运维与应急响应。


一、消息队列基础

1.1 消息队列的作用

作用

说明

异步处理

将非核心流程异步化,提升主流程响应速度

流量削峰

高并发场景下缓冲请求,保护下游服务

系统解耦

生产者与消费者独立演进,降低耦合度

数据管道

在微服务/大数据场景中充当数据传输通道

事件驱动

支撑事件溯源(Event Sourcing)和 CQRS 架构

1.2 典型应用场景

  • 订单系统:下单后异步发送短信、推送、库存扣减

  • 日志收集:应用日志 → MQ → ELK / ClickHouse

  • 实时计算:用户行为 → Kafka → Flink / Spark Streaming

  • 分布式事务:基于 MQ 的最终一致性方案(如 RocketMQ 事务消息)

  • 任务调度:延迟队列实现定时任务(如 RabbitMQ TTL + 死信队列)

1.3 主流消息队列选型对比

特性

RabbitMQ

Apache Kafka

RocketMQ

开发语言

Erlang

Java/Scala

Java

协议

AMQP

自定义协议

自定义协议

消息模型

Exchange + Queue

Topic + Partition

Topic + Queue

吞吐量

万级/s

百万级/s

十万级/s

延迟

微秒级

毫秒级

毫秒级

消息可靠性

高(持久化+确认)

高(副本+ISR)

高(同步刷盘+主从)

事务消息

支持

不支持原生

原生支持

延迟消息

插件支持

不支持原生

原生支持

消息回溯

不支持

支持(按 offset)

支持(按时间戳)

适用场景

业务消息、复杂路由

日志、大数据流

电商、金融交易

运维复杂度

中等

较高

中等

选型建议

  • 业务系统消息路由复杂 → RabbitMQ

  • 日志采集、大数据管道、高吞吐 → Kafka

  • 电商/金融场景、需要事务消息和延迟消息 → RocketMQ


二、RabbitMQ 运维

2.1 安装部署

CentOS / RHEL 安装:

# 安装 Erlang 依赖
curl -s https://packagecloud.io/install/repositories/rabbitmq/erlang/script.rpm.sh | sudo bash
sudo yum install -y erlang

# 安装 RabbitMQ
curl -s https://packagecloud.io/install/repositories/rabbitmq/rabbitmq-server/script.rpm.sh | sudo bash
sudo yum install -y rabbitmq-server

# 启动服务
sudo systemctl enable rabbitmq-server
sudo systemctl start rabbitmq-server

Ubuntu / Debian 安装:

sudo apt-get update
sudo apt-get install -y rabbitmq-server
sudo systemctl enable rabbitmq-server
sudo systemctl start rabbitmq-server

Docker 安装:

docker run -d --name rabbitmq \
  -p 5672:5672 \
  -p 15672:15672 \
  -e RABBITMQ_DEFAULT_USER=admin \
  -e RABBITMQ_DEFAULT_PASS=admin123 \
  -v /data/rabbitmq:/var/lib/rabbitmq \
  rabbitmq:3-management

2.2 管理界面配置

# 启用管理插件
sudo rabbitmq-plugins enable rabbitmq_management

# 创建管理员用户
sudo rabbitmqctl add_user admin admin123
sudo rabbitmqctl set_user_tags admin administrator
sudo rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"

# 删除默认 guest 用户(安全加固)
sudo rabbitmqctl delete_user guest

访问管理界面:http://<host>:15672

关键配置文件 /etc/rabbitmq/rabbitmq.conf

# 监听端口
listeners.tcp.default = 5672
# 管理界面端口
management.tcp.port = 15672
# 内存高水位线(默认 40% 系统内存)
vm_memory_high_watermark.relative = 0.6
# 磁盘空间低水位线
disk_free_limit.absolute = 2GB
# 最大连接数
# channel_max = 2047
# 心跳超时
heartbeat = 60
# 消息持久化
queue_master_locator = min-masters

2.3 Exchange 与 Queue 管理

Exchange 类型:

类型

说明

路由规则

direct

精确匹配

routing_key 完全匹配

topic

模式匹配

支持 *(一个词)和 #(零或多个词)

fanout

广播

忽略 routing_key,投递到所有绑定队列

headers

头部匹配

根据消息 header 匹配

使用 rabbitmqadmin 管理:

# 创建 Exchange
rabbitmqadmin declare exchange name=order.events type=topic durable=true

# 创建 Queue
rabbitmqadmin declare queue name=order.process durable=true

# 绑定 Exchange 到 Queue
rabbitmqadmin declare binding source=order.events destination=order.process routing_key="order.created"

# 发送测试消息
rabbitmqadmin publish exchange=order.events routing_key="order.created" payload='{"orderId":"10001","amount":99.5}'

# 消费消息(获取一条)
rabbitmqadmin get queue=order.process count=1

# 查看队列详情
rabbitmqadmin list queues name messages consumers memory

使用 rabbitmqctl 管理:

# 列出所有队列
sudo rabbitmqctl list_queues name messages consumers memory

# 列出所有 Exchange
sudo rabbitmqctl list_exchanges name type durable

# 列出所有连接
sudo rabbitmqctl list_connections name peer_host state

# 列出所有通道
sudo rabbitmqctl list_channels connection_name number consumer_count

# 清空队列
sudo rabbitmqctl purge_queue order.process

# 关闭连接
sudo rabbitmqctl close_connection "<connection_id>" "admin close"

2.4 集群部署

三节点集群架构:

Node1 (rabbit@node1)  ←→  Node2 (rabbit@node2)  ←→  Node3 (rabbit@node3)
        ↓                        ↓                        ↓
     HAProxy                  HAProxy                  HAProxy
        ↓                        ↓                        ↓
              Client (连接任意节点)

搭建步骤:

# === 所有节点执行 ===
# 1. 同步 Erlang Cookie(关键!所有节点必须一致)
scp node1:/var/lib/rabbitmq/.erlang.cookie /var/lib/rabbitmq/.erlang.cookie
chmod 400 /var/lib/rabbitmq/.erlang.cookie
chown rabbitmq:rabbitmq /var/lib/rabbitmq/.erlang.cookie

# 2. 重启服务
sudo systemctl restart rabbitmq-server

# === 在 node2 和 node3 上执行 ===
# 3. 加入集群
sudo rabbitmqctl stop_app
sudo rabbitmqctl reset
sudo rabbitmqctl join_cluster rabbit@node1
sudo rabbitmqctl start_app

# === 在任意节点验证 ===
sudo rabbitmqctl cluster_status

2.5 镜像队列(高可用)

# 设置镜像队列策略(经典模式 - 3.8 及以下)
sudo rabbitmqctl set_policy ha-all "^" \
  '{"ha-mode":"all","ha-sync-mode":"automatic","ha-promote-on-shutdown":"always"}' \
  --priority 1 --apply-to queues

# 设置镜像队列策略(指定副本数为 3)
sudo rabbitmqctl set_policy ha-three "^order\." \
  '{"ha-mode":"exactly","ha-params":3,"ha-sync-mode":"automatic"}' \
  --priority 1 --apply-to queues

# 3.8+ 推荐使用 Quorum Queues(Raft 协议)
rabbitmqadmin declare queue name=order.process durable=true \
  x-queue-type=quorum arguments='{"x-quorum-initial-group-size":3}'

Quorum Queue vs 经典镜像队列:

特性

经典镜像队列

Quorum Queue

一致性协议

非 Raft

Raft

消息丢失风险

有(脑裂时)

极低

性能

较高

略低

推荐版本

3.8 以下

3.8+ 推荐


三、Apache Kafka 运维

3.1 安装部署

基础环境准备:

# 安装 Java(Kafka 3.x 需要 Java 11+)
sudo yum install -y java-11-openjdk java-11-openjdk-devel

# 设置 JAVA_HOME
echo 'export JAVA_HOME=/usr/lib/jvm/java-11-openjdk' >> /etc/profile
source /etc/profile

安装 Kafka:

# 下载
KAFKA_VERSION="3.7.0"
wget https://downloads.apache.org/kafka/${KAFKA_VERSION}/kafka_2.13-${KAFKA_VERSION}.tgz
tar -xzf kafka_2.13-${KAFKA_VERSION}.tgz -C /opt/
ln -s /opt/kafka_2.13-${KAFKA_VERSION} /opt/kafka

# 创建数据目录
mkdir -p /data/kafka-logs

核心配置 config/server.properties

# Broker 唯一标识
broker.id=0
# 监听地址
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://kafka1:9092
# 日志目录(多个目录可分散 IO)
log.dirs=/data/kafka-logs
# 分区数
num.partitions=3
# 副本数
default.replication.factor=3
# 最小同步副本(用于 acks=all)
min.insync.replicas=2
# 消息保留时间
log.retention.hours=168
# 单个日志段大小
log.segment.bytes=1073741824
# ZooKeeper 连接
zookeeper.connect=zk1:2181,zk2:2181,zk3:2181/kafka
# 自动创建 Topic
auto.create.topics.enable=false
# 消息最大大小
message.max.bytes=10485760
# 副本同步最大大小
replica.fetch.max.bytes=10485760

KRaft 模式(无 ZooKeeper,3.3+ 推荐):

# 生成集群 UUID
KAFKA_CLUSTER_ID=$(kafka-storage.sh random-uuid)

# 格式化存储目录
kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties

# 启动(无需 ZooKeeper)
kafka-server-start.sh -daemon config/kraft/server.properties

KRaft 配置 config/kraft/server.properties

process.roles=broker,controller
node.id=1
controller.quorum.voters=1@kafka1:9093,2@kafka2:9093,3@kafka3:9093
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
advertised.listeners=PLAINTEXT://kafka1:9092
log.dirs=/data/kafka-logs

3.2 Broker 管理

# 启动 Broker
kafka-server-start.sh -daemon config/server.properties

# 停止 Broker
kafka-server-stop.sh

# 查看 Broker 列表
kafka-broker-api-versions.sh --bootstrap-server kafka1:9092

# 查看 Broker 配置
kafka-configs.sh --bootstrap-server kafka1:9092 \
  --entity-type brokers --entity-name 0 --describe

# 动态修改 Broker 配置
kafka-configs.sh --bootstrap-server kafka1:9092 \
  --entity-type brokers --entity-name 0 \
  --alter --add-config log.retention.hours=240

3.3 Topic 管理

# 创建 Topic
kafka-topics.sh --bootstrap-server kafka1:9092 --create \
  --topic order-events \
  --partitions 6 \
  --replication-factor 3 \
  --config retention.ms=604800000 \
  --config max.message.bytes=10485760

# 列出所有 Topic
kafka-topics.sh --bootstrap-server kafka1:9092 --list

# 查看 Topic 详情
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --describe --topic order-events

# 增加分区(只能增不能减)
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --alter --topic order-events --partitions 12

# 修改 Topic 配置
kafka-configs.sh --bootstrap-server kafka1:9092 \
  --entity-type topics --entity-name order-events \
  --alter --add-config retention.ms=259200000

# 删除 Topic(需确认 delete.topic.enable=true)
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --delete --topic order-events

3.4 Consumer Group 管理

# 列出所有消费者组
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 --list

# 查看消费者组详情(含 lag)
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --describe --group order-processor

# 输出示例:
# GROUP           TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
# order-processor order-events    0          150000          155000          5000
# order-processor order-events    1          200000          201000          1000

# 重置消费者 offset(到最早位置)
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --group order-processor --topic order-events \
  --reset-offsets --to-earliest --execute

# 重置到指定时间点
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --group order-processor --topic order-events \
  --reset-offsets --to-datetime "2026-06-01T00:00:00.000" --execute

# 删除消费者组
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --delete --group order-processor

3.5 分区与副本管理

# 查看 Topic 分区分配
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --describe --topic order-events --under-replicated-partitions

# 分区重分配(扩缩容后迁移数据)
# 1. 生成重分配计划
kafka-reassign-partitions.sh --bootstrap-server kafka1:9092 \
  --generate \
  --topics-to-move-json-file topics.json \
  --broker-list "0,1,2,3"

# 2. 执行重分配
kafka-reassign-partitions.sh --bootstrap-server kafka1:9092 \
  --execute \
  --reassignment-json-file reassignment.json

# 3. 验证重分配进度
kafka-reassign-partitions.sh --bootstrap-server kafka1:9092 \
  --verify \
  --reassignment-json-file reassignment.json

# 修改副本因子(增加副本)
cat > increase-replication.json << EOF
{
  "version": 1,
  "partitions": [
    {"topic": "order-events", "partition": 0, "replicas": [0,1,2]},
    {"topic": "order-events", "partition": 1, "replicas": [1,2,0]},
    {"topic": "order-events", "partition": 2, "replicas": [2,0,1]}
  ]
}
EOF

kafka-reassign-partitions.sh --bootstrap-server kafka1:9092 \
  --execute --reassignment-json-file increase-replication.json

3.6 Kafka Connect

# 启动分布式 Connect Worker
connect-distributed.sh -daemon config/connect-distributed.properties

# 创建 JDBC Source Connector
curl -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "mysql-source",
    "config": {
      "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
      "connection.url": "jdbc:mysql://db-host:3306/mydb",
      "connection.user": "reader",
      "connection.password": "secret",
      "table.whitelist": "orders,users",
      "mode": "incrementing",
      "incrementing.column.name": "id",
      "topic.prefix": "cdc-"
    }
  }'

# 创建 File Sink Connector
curl -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "file-sink",
    "config": {
      "connector.class": "FileStreamSink",
      "tasks.max": 1,
      "file": "/data/sink/output.txt",
      "topics": "order-events"
    }
  }'

# 查看 Connector 状态
curl -s http://localhost:8083/connectors/mysql-source/status | jq

# 删除 Connector
curl -X DELETE http://localhost:8083/connectors/mysql-source

3.7 Kafka 集群运维

# 查看集群信息
kafka-metadata.sh --snapshot /data/kafka-logs/__cluster_metadata-0/00000000000000000000.log --cluster-id <id>

# Preferred Leader 选举(自动均衡 Leader)
kafka-leader-election.sh --bootstrap-server kafka1:9092 \
  --election-type preferred \
  --topic order-events --partition 0

# Unclean Leader 选举(数据可能丢失,紧急使用)
kafka-leader-election.sh --bootstrap-server kafka1:9092 \
  --election-type unclean \
  --topic order-events --partition 0

# 查看 Topic 的 Log 详情
kafka-dump-log.sh --files /data/kafka-logs/order-events-0/00000000000000000000.log \
  --print-data-log

四、RocketMQ 运维

4.1 安装部署

# 安装 Java 8+
sudo yum install -y java-1.8.0-openjdk java-1.8.0-openjdk-devel

# 下载 RocketMQ
ROCKETMQ_VERSION="5.2.0"
wget https://downloads.apache.org/rocketmq/${ROCKETMQ_VERSION}/rocketmq-all-${ROCKETMQ_VERSION}-bin-release.zip
unzip rocketmq-all-${ROCKETMQ_VERSION}-bin-release.zip -d /opt/
ln -s /opt/rocketmq-all-${ROCKETMQ_VERSION}-bin-release /opt/rocketmq

# 创建数据目录
mkdir -p /data/rocketmq/store /data/rocketmq/logs

# 环境变量
export ROCKETMQ_HOME=/opt/rocketmq
export PATH=$ROCKETMQ_HOME/bin:$PATH

4.2 NameServer 部署

# 修改 JVM 参数(namesrv.sh 或 runserver.sh)
# 生产环境建议:
# -Xms4g -Xmx4g -Xmn2g

# 启动 NameServer
nohup sh mqnamesrv &

# 验证 NameServer 启动
tail -f ~/logs/rocketmqlogs/namesrv.log
# 看到 "The Name Server boot success" 表示成功

# 后台启动(推荐)
nohup sh mqnamesrv -n "namesrv1:9876;namesrv2:9876" > /dev/null 2>&1 &

4.3 Broker 部署

Broker 配置文件 conf/broker.conf

# Broker 名称(同一主从需相同)
brokerName=broker-a
# Broker ID(0=Master, >0=Slave)
brokerId=0
# NameServer 地址
namesrvAddr=namesrv1:9876;namesrv2:9876
# 存储路径
storePathRootDir=/data/rocketmq/store
storePathCommitLog=/data/rocketmq/store/commitlog
# 刷盘方式:ASYNC_FLUSH(异步)/ SYNC_FLUSH(同步)
flushDiskType=ASYNC_FLUSH
# 同步复制/异步复制
brokerRole=SYNC_MASTER
# 自动创建 Topic
autoCreateTopicEnable=false
# 消费者组允许从 slave 消费
slaveReadEnable=true
# 文件保留时间(小时)
fileReservedTime=72
# 单个 CommitLog 文件大小
mapedFileSizeCommitLog=1073741824
# 单个 ConsumeQueue 文件存储条数
mapedFileSizeConsumeQueue=300000
# 发送消息线程池
sendMessageThreadPoolNums=64
# 消费消息线程池
consumeMessageThreadPoolNums=64
# 启动 Master Broker
nohup sh mqbroker -c conf/broker.conf > /dev/null 2>&1 &

# 启动 Slave Broker(修改 brokerId=1, brokerRole=SLAVE)
nohup sh mqbroker -c conf/broker-slave.conf > /dev/null 2>&1 &

# 验证 Broker 注册到 NameServer
sh mqadmin clusterList -n namesrv1:9876

4.4 事务消息

# 生产者发送事务消息(Java 示例核心逻辑)
# 1. 发送 Half 消息
# 2. 执行本地事务
# 3. 根据结果提交或回滚
# 4. 如果无响应,Broker 回查本地事务状态

# 查看事务消息状态
sh mqadmin queryMsgById -n namesrv1:9876 -i <msgId>

# 查看 Half 消息(存储在 RMQ_SYS_TRANS_HALF_TOPIC)
sh mqadmin topicStatus -n namesrv1:9876 -t RMQ_SYS_TRANS_HALF_TOPIC

4.5 RocketMQ 集群模式

模式

说明

适用场景

单 Master

单节点

开发/测试

多 Master

多节点无副本

允许少量丢失

多 Master 多 Slave(异步)

主从异步复制

高吞吐

多 Master 多 Slave(同步)

主从同步复制

高可靠(推荐)

DLedger

Raft 协议自动选主

金融级高可用

DLedger 模式配置(3 副本自动选主):

# broker-dledger.conf
brokerName=broker-a
enableDLegerCommitLog=true
dLegerGroup=broker-a
dLegerPeers=n0-kafka1:40911;n1-kafka2:40912;n2-kafka3:40913
dLegerSelfId=n0
namesrvAddr=namesrv1:9876;namesrv2:9876
# 启动 DLedger Broker
nohup sh mqbroker -c conf/broker-dledger.conf > /dev/null 2>&1 &

# 查看 DLedger 状态
sh mqadmin getBrokerConfig -n namesrv1:9876 -b kafka1:10911

五、消息可靠性保障

5.1 生产端确认

RabbitMQ Publisher Confirm:

# Python pika 示例
channel.confirm_delivery()  # 开启确认模式

try:
    channel.basic_publish(
        exchange='order.events',
        routing_key='order.created',
        body=json.dumps(order),
        properties=pika.BasicProperties(
            delivery_mode=2,  # 持久化消息
            content_type='application/json'
        ),
        mandatory=True  # 无法路由时返回
    )
    print("消息发送成功")
except pika.exceptions.UnroutableError:
    print("消息无法路由")
except pika.exceptions.NackError:
    print("Broker 拒绝消息")

Kafka Producer 确认:

# Producer 配置
acks=all                    # 所有 ISR 副本确认
retries=2147483647          # 无限重试
enable.idempotence=true     # 幂等生产者
max.in.flight.requests.per.connection=5

RocketMQ 同步发送确认:

DefaultMQProducer producer = new DefaultMQProducer("producer-group");
producer.setNamesrvAddr("namesrv1:9876");
producer.setRetryTimesWhenSendFailed(3);      // 同步重试次数
producer.setRetryTimesWhenSendAsyncFailed(3); // 异步重试次数
producer.setSendMsgTimeout(3000);              // 超时时间
producer.start();

SendResult result = producer.send(message);
// result.getSendStatus() == SendStatus.SEND_OK

5.2 消费端确认

RabbitMQ 手动 ACK:

def callback(ch, method, properties, body):
    try:
        process_message(body)
        ch.basic_ack(delivery_tag=method.delivery_tag)  # 确认
    except Exception as e:
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=False  # 不重新入队,进入死信队列
        )

channel.basic_consume(queue='order.process', on_message_callback=callback)
channel.basic_qos(prefetch_count=100)  # 预取数量

Kafka Consumer 手动提交:

enable.auto.commit=false
auto.offset.reset=earliest
consumer.subscribe(Arrays.asList("order-events"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        try {
            processRecord(record);
        } catch (Exception e) {
            handleFailure(record, e);
        }
    }
    consumer.commitSync();  // 手动同步提交
}

RocketMQ 并发消费确认:

DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer-group");
consumer.setNamesrvAddr("namesrv1:9876");
consumer.subscribe("order-events", "*");
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
    try {
        for (MessageExt msg : msgs) {
            processMessage(msg);
        }
        return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
    } catch (Exception e) {
        return ConsumeConcurrentlyStatus.RECONSUME_LATER;  // 稍后重试
    }
});
consumer.start();

5.3 消费幂等

// 基于消息 ID 的幂等方案
public boolean processIdempotent(MessageExt msg) {
    String msgId = msg.getMsgId();
    // 1. 查询 Redis 是否已处理
    Boolean exists = redisTemplate.hasKey("msg:processed:" + msgId);
    if (Boolean.TRUE.equals(exists)) {
        log.info("消息已处理,跳过: {}", msgId);
        return true;
    }
    // 2. 执行业务逻辑
    boolean success = doBusiness(msg);
    // 3. 标记已处理(设置过期时间)
    if (success) {
        redisTemplate.opsForValue().set(
            "msg:processed:" + msgId, "1", 24, TimeUnit.HOURS
        );
    }
    return success;
}

5.4 死信队列(DLQ)

RabbitMQ 死信队列配置:

# 创建死信队列
rabbitmqadmin declare queue name=order.dead-letter durable=true

# 创建原始队列时指定死信 Exchange
rabbitmqadmin declare queue name=order.process durable=true \
  arguments='{"x-dead-letter-exchange":"","x-dead-letter-routing-key":"order.dead-letter","x-message-ttl":60000,"x-max-length":10000}'

# 绑定死信队列
rabbitmqadmin declare binding source="" destination=order.dead-letter routing_key="order.dead-letter"

RocketMQ 死信队列:

# 查看死信队列(自动创建,格式:%DLQ%消费者组名)
sh mqadmin topicList -n namesrv1:9876 | grep DLQ

# 查看死信消息
sh mqadmin queryMsgByKey -n namesrv1:9876 -t "%DLQ%consumer-group" -k <messageKey>

# 消费死信消息(重新投递到原始 Topic)
sh mqadmin consumeMessage -n namesrv1:9876 -t "%DLQ%consumer-group" -o <offset>

5.5 重试机制

RabbitMQ 重试(基于插件或手动实现):

# 基于 x-death 头的重试计数
def callback(ch, method, properties, body):
    retry_count = 0
    if properties.headers and 'x-death' in properties.headers:
        retry_count = properties.headers['x-death'][0].get('count', 0)

    if retry_count >= 3:
        # 超过重试次数,发送到死信队列
        ch.basic_ack(delivery_tag=method.delivery_tag)
        send_to_dlq(body)
        return

    try:
        process_message(body)
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except Exception:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

RocketMQ 延迟重试等级:

重试次数

延迟时间

1

10s

2

30s

3

1min

4

2min

5

3min

6

4min

7

5min

8

6min

9

7min

10

8min

11

9min

12

10min

13

20min

14

30min

15

1h

16

2h

# 查看重试队列消息数
sh mqadmin topicStatus -n namesrv1:9876 -t "%RETRY%consumer-group"

六、性能优化

6.1 批量发送

Kafka Producer 批量配置:

batch.size=16384              # 批次大小(字节)
linger.ms=5                   # 等待时间(毫秒)
buffer.memory=33554432        # 缓冲区大小(32MB)
compression.type=lz4          # 压缩类型

RocketMQ 批量发送:

List<Message> messages = new ArrayList<>();
for (int i = 0; i < 100; i++) {
    messages.add(new Message("order-events", "TAG_A", ("msg-" + i).getBytes()));
}
SendResult result = producer.send(messages);  // 批量发送,单次最多 32KB

6.2 消息压缩

压缩算法

压缩比

CPU 开销

推荐场景

none

小消息

gzip

存储优化

snappy

Kafka 默认推荐

lz4

高吞吐场景

zstd

Kafka 2.1+ 推荐

6.3 分区优化

Kafka 分区数计算公式:

目标分区数 = max(目标吞吐量 / 单分区吞吐量, 消费者数量)

# 示例:目标 100MB/s,单分区 30MB/s,消费者 8 个
# 分区数 = max(100/30, 8) = max(3.3, 8) = 8(取 12 个留余量)

RabbitMQ 队列优化:

# 惰性队列(大量消息时使用磁盘存储,降低内存占用)
rabbitmqadmin declare queue name=large.queue durable=true \
  arguments='{"x-queue-mode":"lazy"}'

6.4 消息积压处理

Kafka 消费积压排查:

# 1. 查看消费者 lag
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --describe --group order-processor

# 2. 临时增加消费者实例(不超过分区数)
# 3. 如果分区数不够,先扩分区
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --alter --topic order-events --partitions 24

# 4. 紧急方案:重置 offset 跳过积压消息
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --group order-processor --topic order-events \
  --reset-offsets --to-latest --execute

RabbitMQ 消费积压排查:

# 查看队列消息数
rabbitmqadmin list queues name messages consumers memory

# 查看消费者详情
rabbitmqctl list_consumers

# 紧急清空队列
sudo rabbitmqctl purge_queue order.process

# 增加消费者 prefetch(提升并发)
# channel.basic_qos(prefetch_count=500)

七、监控与告警

7.1 RabbitMQ 监控

# 启用 Prometheus 插件
sudo rabbitmq-plugins enable rabbitmq_prometheus

# Prometheus 配置
cat >> prometheus.yml << EOF
scrape_configs:
  - job_name: 'rabbitmq'
    static_configs:
      - targets: ['rabbitmq1:15692']
    metrics_path: /metrics
EOF

关键监控指标:

指标

含义

告警阈值

rabbitmq_queue_messages

队列消息总数

> 10000

rabbitmq_queue_messages_ready

待消费消息数

> 5000

rabbitmq_queue_messages_unacked

未确认消息数

> 1000

rabbitmq_queue_consumers

消费者数量

< 1

rabbitmq_channel_messages_published_total

发布消息速率

突降 50%

rabbitmq_process_resident_memory_bytes

内存使用

> 80% 水位线

rabbitmq_disk_space_available_bytes

可用磁盘

< 2GB

7.2 Kafka 监控

# JMX 监控(启动时开启)
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote \
  -Dcom.sun.management.jmxremote.port=9999 \
  -Dcom.sun.management.jmxremote.authenticate=false \
  -Dcom.sun.management.jmxremote.ssl=false"

# 使用 kafka-exporter(Prometheus 方案)
docker run -d --name kafka-exporter \
  -p 9308:9308 \
  danielqsj/kafka-exporter \
  --kafka.server=kafka1:9092

关键监控指标:

指标

含义

告警阈值

kafka_server_BrokerTopicMetrics_MessagesInPerSec

消息入站速率

突降 50%

kafka_server_ReplicaManager_UnderReplicatedPartitions

副本不足分区数

> 0

kafka_consumer_group_lag

消费者 Lag

> 100000

kafka_server_BrokerTopicMetrics_BytesInPerSec

字节入站速率

监控趋势

kafka_server_ReplicaManager_OfflineReplicaCount

离线副本数

> 0

kafka_server_BrokerTopicMetrics_TotalProduceRequestsPerSec

生产请求速率

突降

7.3 RocketMQ 监控

# 使用 RocketMQ Dashboard
docker run -d --name rocketmq-dashboard \
  -p 8080:8080 \
  -e "JAVA_OPTS=-Drocketmq.namesrv.addr=namesrv1:9876" \
  apache/rocketmq-dashboard:latest

# 命令行监控
# 查看集群信息
sh mqadmin clusterList -n namesrv1:9876

# 查看 Topic 列表及消息统计
sh mqadmin topicStatus -n namesrv1:9876

# 查看消费者进度
sh mqadmin consumerProgress -n namesrv1:9876 -g consumer-group

# 查看 Broker 状态
sh mqadmin brokerStatus -n namesrv1:9876 -b kafka1:10911

7.4 Grafana Dashboard 推荐

MQ

Dashboard ID

说明

RabbitMQ

10991

RabbitMQ Overview

Kafka

7589

Kafka Overview

Kafka Lag

12693

Kafka Consumer Lag

RocketMQ

10478

RocketMQ Exporter


八、常见问题排查

8.1 RabbitMQ 常见问题

问题 1:队列阻塞(flow / alarm 状态)

# 查看是否有 alarm
sudo rabbitmqctl status | grep -A 5 alarms

# 常见原因:内存超水位线
# 解决方案:
# 1. 增加内存水位线
sudo rabbitmqctl set_vm_memory_high_watermark 0.7

# 2. 清理不需要的队列
rabbitmqadmin list queues name messages memory | sort -k3 -rn | head -20

# 3. 扩容节点

问题 2:连接数过多

# 查看连接状态
sudo rabbitmqctl list_connections name peer_host state channels

# 关闭空闲连接
sudo rabbitmqctl list_connections name state idle_since | grep idle

# 增加最大连接限制(rabbitmq.conf)
# channel_max = 4096

问题 3:消息堆积导致内存 OOM

# 切换到惰性队列模式(运行时)
rabbitmqadmin declare queue name=order.process durable=true \
  arguments='{"x-queue-mode":"lazy"}'

# 或设置队列最大长度
rabbitmqadmin declare queue name=order.process durable=true \
  arguments='{"x-max-length":1000000,"x-overflow":"drop-head"}'

8.2 Kafka 常见问题

问题 1:分区 Leader 选举失败

# 查看离线分区
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --describe --under-replicated-partitions

# 手动触发 Leader 选举
kafka-leader-election.sh --bootstrap-server kafka1:9092 \
  --election-type preferred \
  --topic order-events --partition 0

# 如果所有副本都不可用(紧急)
kafka-leader-election.sh --bootstrap-server kafka1:9092 \
  --election-type unclean \
  --topic order-events --partition 0

问题 2:消费 Lag 持续增长

# 1. 检查消费者是否存活
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --describe --group order-processor --state

# 2. 检查消费者是否有异常日志
# 3. 增加消费者实例数(不超过分区数)
# 4. 优化消费逻辑(减少处理时间)
# 5. 临时方案:重置 offset

问题 3:磁盘空间不足

# 查看 Topic 大小
du -sh /data/kafka-logs/*

# 紧急清理:缩短保留时间
kafka-configs.sh --bootstrap-server kafka1:9092 \
  --entity-type topics --entity-name order-events \
  --alter --add-config retention.ms=86400000

# 或删除不需要的 Topic
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --delete --topic old-topic

# 手动触发日志清理
kafka-delete-records.sh --bootstrap-server kafka1:9092 \
  --offset-json-file delete-records.json

问题 4:消息丢失排查

# 检查 Producer 配置
# acks=all, retries>0, enable.idempotence=true

# 检查 Broker 配置
# min.insync.replicas=2 (不能大于 replication.factor)
# unclean.leader.election.enable=false

# 检查 Consumer 配置
# enable.auto.commit=false, 手动提交 offset

8.3 RocketMQ 常见问题

问题 1:消息发送超时

# 检查 NameServer 是否正常
sh mqadmin clusterList -n namesrv1:9876

# 检查 Broker 是否存活
sh mqadmin brokerStatus -n namesrv1:9876 -b broker-host:10911

# 检查网络连通性
telnet broker-host 10911

# 查看 Broker 日志
tail -f ~/logs/rocketmqlogs/broker.log | grep -i "slow\|timeout\|error"

问题 2:消费失败不断重试

# 查看重试队列消息数
sh mqadmin topicStatus -n namesrv1:9876 -t "%RETRY%consumer-group"

# 查看具体消息
sh mqadmin queryMsgByKey -n namesrv1:9876 \
  -t "%RETRY%consumer-group" -k <messageKey>

# 跳过重试消息(重置 offset)
sh mqadmin resetOffsetByTime -n namesrv1:9876 \
  -g consumer-group -t order-events -s now

问题 3:CommitLog 写入缓慢

# 检查磁盘 IO
iostat -x 1 10

# 查看刷盘耗时
grep "putMessage" ~/logs/rocketmqlogs/store.log | tail -20

# 优化方案:
# 1. 使用 SSD 磁盘
# 2. 改为异步刷盘(flushDiskType=ASYNC_FLUSH)
# 3. 增加 CommitLog 文件预分配

九、运维命令速查表

9.1 RabbitMQ 速查

操作

命令

查看服务状态

sudo rabbitmqctl status

查看集群状态

sudo rabbitmqctl cluster_status

列出所有队列

sudo rabbitmqctl list_queues name messages consumers

列出所有连接

sudo rabbitmqctl list_connections

列出所有通道

sudo rabbitmqctl list_channels

添加用户

sudo rabbitmqctl add_user <user> <pass>

设置权限

sudo rabbitmqctl set_permissions -p / <user> ".*" ".*" ".*"

清空队列

sudo rabbitmqctl purge_queue <queue>

关闭连接

sudo rabbitmqctl close_connection <id> "reason"

设置策略

sudo rabbitmqctl set_policy <name> <pattern> <json>

移除策略

sudo rabbitmqctl clear_policy <name>

启用插件

sudo rabbitmq-plugins enable <plugin>

列出插件

sudo rabbitmq-plugins list

加入集群

sudo rabbitmqctl join_cluster rabbit@<node>

移除节点

sudo rabbitmqctl forget_cluster_node rabbit@<node>

重置节点

sudo rabbitmqctl reset

9.2 Kafka 速查

操作

命令

创建 Topic

kafka-topics.sh --bootstrap-server <addr> --create --topic <t> --partitions <n> --replication-factor <n>

列出 Topic

kafka-topics.sh --bootstrap-server <addr> --list

查看 Topic

kafka-topics.sh --bootstrap-server <addr> --describe --topic <t>

删除 Topic

kafka-topics.sh --bootstrap-server <addr> --delete --topic <t>

增加分区

kafka-topics.sh --bootstrap-server <addr> --alter --topic <t> --partitions <n>

发送消息

kafka-console-producer.sh --bootstrap-server <addr> --topic <t>

消费消息

kafka-console-consumer.sh --bootstrap-server <addr> --topic <t> --from-beginning

查看消费者组

kafka-consumer-groups.sh --bootstrap-server <addr> --list

查看消费 Lag

kafka-consumer-groups.sh --bootstrap-server <addr> --describe --group <g>

重置 Offset

kafka-consumer-groups.sh --bootstrap-server <addr> --group <g> --topic <t> --reset-offsets --to-earliest --execute

Leader 选举

kafka-leader-election.sh --bootstrap-server <addr> --election-type preferred --topic <t> --partition <p>

分区重分配

kafka-reassign-partitions.sh --bootstrap-server <addr> --execute --reassignment-json-file <f>

查看 Broker 配置

kafka-configs.sh --bootstrap-server <addr> --entity-type brokers --entity-name <id> --describe

9.3 RocketMQ 速查

操作

命令

启动 NameServer

mqnamesrv

启动 Broker

mqbroker -c broker.conf

查看集群

mqadmin clusterList -n <ns>

查看 Topic 列表

mqadmin topicList -n <ns>

创建 Topic

mqadmin updateTopic -n <ns> -b <broker> -t <topic>

删除 Topic

mqadmin deleteTopic -n <ns> -b <broker> -t <topic>

查看 Topic 状态

mqadmin topicStatus -n <ns> -t <topic>

查看消费者进度

mqadmin consumerProgress -n <ns> -g <group>

查看 Broker 状态

mqadmin brokerStatus -n <ns> -b <broker>

发送测试消息

mqadmin sendMsgStatus -n <ns> -t <topic>

消费测试消息

mqadmin consumeMessage -n <ns> -t <topic>

查询消息(按 ID)

mqadmin queryMsgById -n <ns> -i <msgId>

查询消息(按 Key)

mqadmin queryMsgByKey -n <ns> -t <topic> -k <key>

重置消费位点

mqadmin resetOffsetByTime -n <ns> -g <group> -t <topic> -s <timestamp>

更新 Broker 配置

mqadmin updateBrokerConfig -n <ns> -b <broker> -k <key> -v <value>

获取 Broker 配置

mqadmin getBrokerConfig -n <ns> -b <broker>


十、运维最佳实践

10.1 容量规划

  • 磁盘:根据消息量 × 消息大小 × 保留天数 × 副本数计算

  • 内存:RabbitMQ 注意内存水位线;Kafka 主要用 Page Cache

  • 网络:跨机房部署时考虑带宽成本和延迟

  • CPU:压缩和加密会显著增加 CPU 开销

10.2 安全加固

# RabbitMQ:启用 TLS
listeners.ssl.default = 5671
ssl_options.cacertfile = /etc/rabbitmq/ca_certificate.pem
ssl_options.certfile = /etc/rabbitmq/server_certificate.pem
ssl_options.keyfile = /etc/rabbitmq/server_key.pem
ssl_options.verify = verify_peer
ssl_options.fail_if_no_peer_cert = true

# Kafka:启用 SASL + SSL
security.protocol=SASL_SSL
sasl.mechanism=SCRAM-SHA-256
ssl.truststore.location=/etc/kafka/truststore.jks
ssl.keystore.location=/etc/kafka/keystore.jks

# RocketMQ:启用 ACL
aclEnable=true
# 配置 plain_acl.yml 定义权限

10.3 备份与恢复

# RabbitMQ 元数据导出
sudo rabbitmqctl export_definitions /backup/rabbitmq-defs.json

# RabbitMQ 元数据导入
sudo rabbitmqctl import_definitions /backup/rabbitmq-defs.json

# Kafka 数据备份(使用 kafka-reassign-partitions 将副本迁移到备份集群)
# RocketMQ CommitLog 目录备份
rsync -avz /data/rocketmq/store/ backup-host:/backup/rocketmq-store/

10.4 版本升级注意

MQ

升级建议

RabbitMQ

滚动升级,先升级从节点再主节点,注意 Erlang 版本兼容

Kafka

滚动重启,注意 Inter-Broker Protocol Version 配置

RocketMQ

先升级 NameServer 再 Broker,注意消息格式兼容


版本信息:本文基于 RabbitMQ 3.12+、Apache Kafka 3.7+、RocketMQ 5.2+ 编写。不同版本的配置项和命令可能存在差异,请参考对应版本的官方文档。