本文系统梳理消息队列(MQ)的运维知识体系,覆盖 RabbitMQ、Kafka、RocketMQ 三大主流中间件的安装部署、集群管理、消息可靠性保障、性能优化与故障排查,适用于生产环境的日常运维与应急响应。
一、消息队列基础
1.1 消息队列的作用
1.2 典型应用场景
订单系统:下单后异步发送短信、推送、库存扣减
日志收集:应用日志 → MQ → ELK / ClickHouse
实时计算:用户行为 → Kafka → Flink / Spark Streaming
分布式事务:基于 MQ 的最终一致性方案(如 RocketMQ 事务消息)
任务调度:延迟队列实现定时任务(如 RabbitMQ TTL + 死信队列)
1.3 主流消息队列选型对比
选型建议:
业务系统消息路由复杂 → 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-serverUbuntu / Debian 安装:
sudo apt-get update
sudo apt-get install -y rabbitmq-server
sudo systemctl enable rabbitmq-server
sudo systemctl start rabbitmq-serverDocker 安装:
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-management2.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-masters2.3 Exchange 与 Queue 管理
Exchange 类型:
使用 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_status2.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 经典镜像队列:
三、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=10485760KRaft 模式(无 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.propertiesKRaft 配置 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-logs3.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=2403.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-events3.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-processor3.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.json3.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-source3.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:$PATH4.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:98764.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_TOPIC4.5 RocketMQ 集群模式
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=5RocketMQ 同步发送确认:
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_OK5.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=earliestconsumer.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 延迟重试等级:
# 查看重试队列消息数
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); // 批量发送,单次最多 32KB6.2 消息压缩
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 --executeRabbitMQ 消费积压排查:
# 查看队列消息数
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关键监控指标:
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关键监控指标:
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:109117.4 Grafana Dashboard 推荐
八、常见问题排查
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, 手动提交 offset8.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 速查
9.2 Kafka 速查
9.3 RocketMQ 速查
十、运维最佳实践
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 版本升级注意
版本信息:本文基于 RabbitMQ 3.12+、Apache Kafka 3.7+、RocketMQ 5.2+ 编写。不同版本的配置项和命令可能存在差异,请参考对应版本的官方文档。