一、为什么要做RocketMQ多集群容灾?
先给大家说个真实的坑:去年我帮一家做生鲜配送的公司排查故障,他们的订单系统全靠RocketMQ存订单消息,结果某天机房空调坏了,整个RocketMQ集群直接挂了。更坑的是,他们当时只有这一个集群,挂了之后订单全乱了——有的用户付了钱没收到确认,有的骑手拿不到配送单,整整折腾了4个小时才恢复,赔了客户几十万。
这就是单集群的致命问题:只要集群出问题,整个依赖它的业务都得停。多集群容灾的核心,就是不让一个集群的故障波及所有业务,而且故障发生时,能快速切换到正常的集群,尽量不影响用户体验。
二、核心思路:VIP切换+SDK智能感知
很多人做容灾的第一反应是“集群挂了就切IP”,但直接改IP有两个大问题:一是业务代码要手动改配置重启,慢得要死;二是如果集群没完全挂,只是部分节点出问题,乱切换反而会出乱子。
我们的方案是把两个核心技术结合起来:
- VIP切换:给两个RocketMQ集群绑同一个虚拟IP(VIP),平时只有主集群用这个IP,主集群出问题时,自动把VIP切到备集群;
- SDK智能感知:业务用的RocketMQ客户端(SDK),能自己感知到集群状态,自动选正常的集群发消息,不用人工干预。
这么做的好处是:对业务代码完全透明——代码里写的永远是那个VIP,不用管后面的集群怎么换;切换速度快,能控制在几秒内完成;而且不会乱切,只有真的出问题才换。
三、具体怎么落地?
3.1 第一步:搭两个完全一样的RocketMQ集群
首先得有两个一模一样的集群,不然切过去也用不了。比如主集群叫Cluster-A,备集群叫Cluster-B,每个集群都有3个NameServer、5个Broker(主从同步),配置完全一样,比如队列数、主题、权限都得对齐。
举个具体的配置例子,先给两个集群都配置NameServer的地址(用JSON格式,统一用Java生态的RocketMQ客户端,所以配置都是Java用的):
{
"clusterA": {
"nameServer": "ns-a1:9876,ns-a2:9876,ns-a3:9876"
},
"clusterB": {
"nameServer": "ns-b1:9876,ns-b2:9876,ns-b3:9876"
}
}
这里要注意两个细节:一是两个集群的主题必须完全一致,比如订单主题叫order-topic,两个集群都得有这个主题,而且分区数一样;二是备集群的数据要和主集群同步,平时备集群是只读的,只有切换过去才变成可写。
3.2 第二步:给两个集群绑VIP,实现自动切换
VIP其实就是一个固定的IP地址,比如我们选192.168.1.100作为这个业务的RocketMQ服务IP。平时这个IP绑定在Cluster-A的NameServer上,业务通过这个IP访问主集群。
怎么实现自动切换?我们可以用Keepalived这个工具,它能监控集群的状态,出问题时自动把VIP切到备集群。
给Keepalived写个配置文件(用Shell脚本启动,配置用conf格式):
# 主集群的Keepalived配置(Cluster-A用)
global_defs {
notification_email {
admin@company.com
}
notification_email_from keepalived@company.com
smtp_server smtp.company.com
smtp_connect_timeout 30
router_id LVS_DEVEL
}
vrrp_instance VI_1 {
state MASTER # 主集群的状态是MASTER
interface eth0 # 绑定的网卡,根据自己的服务器改
virtual_router_id 51 # 虚拟路由ID,两个集群必须一样
priority 100 # 主集群的优先级更高,平时用主集群
advert_int 1 # 每隔1秒发一次心跳
authentication {
auth_type PASS
auth_pass 1111 # 两个集群的密码必须一样
}
virtual_ipaddress {
192.168.1.100/24 # 我们的VIP
}
}
# 备集群的Keepalived配置(Cluster-B用)
global_defs {
notification_email {
admin@company.com
}
notification_email_from keepalived@company.com
smtp_server smtp.company.com
smtp_connect_timeout 30
router_id LVS_DEVEL
}
vrrp_instance VI_1 {
state BACKUP # 备集群的状态是BACKUP
interface eth0
virtual_router_id 51 # 和主集群一样
priority 90 # 优先级比主集群低,平时不用
advert_int 1
authentication {
auth_type PASS
auth_pass 1111
}
virtual_ipaddress {
192.168.1.100/24
}
}
这个配置的逻辑很简单:主集群优先级100,备集群90,所以平时VIP在主集群;如果主集群挂了,心跳断了,备集群就会自动把VIP抢过去,整个过程不到3秒。
3.3 第三步:改SDK,实现智能感知
现在业务代码里写的是192.168.1.100这个VIP,那SDK怎么知道这个VIP对应的集群是正常的?我们需要给SDK加个监控逻辑,让它能检测集群的健康状态。
举个Java版的SDK扩展例子(统一用Java技术栈,因为RocketMQ最常用的就是Java客户端):
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.common.message.Message;
public class SmartMQProducer {
// 我们的VIP地址,业务代码只需要写这个
private static final String VIP = "192.168.1.100:9876";
// 两个集群的NameServer地址,SDK内部用
private static final String CLUSTER_A = "ns-a1:9876,ns-a2:9876,ns-a3:9876";
private static final String CLUSTER_B = "ns-b1:9876,ns-b2:9876,ns-b3:9876";
// 当前用的集群,默认是主集群
private String currentCluster = CLUSTER_A;
private DefaultMQProducer producer;
public SmartMQProducer() throws Exception {
// 初始化生产者,用当前的集群
producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr(currentCluster);
producer.start();
// 启动一个线程,每隔1秒检查集群状态
new Thread(this::checkClusterHealth).start();
}
// 检查集群健康的逻辑
private void checkClusterHealth() {
while (true) {
try {
// 先检查当前用的集群是不是正常
boolean currentHealth = checkSingleCluster(currentCluster);
if (!currentHealth) {
// 当前集群挂了,检查备集群是不是正常
boolean backupHealth = checkSingleCluster(currentCluster.equals(CLUSTER_A) ? CLUSTER_B : CLUSTER_A);
if (backupHealth) {
// 备集群正常,切换过去
producer.shutdown();
currentCluster = currentCluster.equals(CLUSTER_A) ? CLUSTER_B : CLUSTER_A;
producer = new DefaultMQProducer("order-producer-group");
producer.setNamesrvAddr(currentCluster);
producer.start();
System.out.println("集群切换成功,当前用:" + currentCluster);
} else {
// 两个集群都挂了,报警
System.out.println("两个集群都故障!请紧急处理!");
}
}
Thread.sleep(1000); // 每隔1秒检查一次
} catch (Exception e) {
e.printStackTrace();
}
}
}
// 检查单个集群是不是正常的方法
private boolean checkSingleCluster(String namesrvAddr) {
try {
// 发一个测试消息,能正常发就说明集群正常
DefaultMQProducer testProducer = new DefaultMQProducer("test-group");
testProducer.setNamesrvAddr(namesrvAddr);
testProducer.start();
Message testMsg = new Message("test-topic", "test-body".getBytes());
testProducer.send(testMsg);
testProducer.shutdown();
return true;
} catch (Exception e) {
return false;
}
}
// 业务发消息的方法,和原来的用法完全一样
public void sendMessage(String topic, String body) throws Exception {
Message msg = new Message(topic, body.getBytes());
producer.send(msg);
}
public static void main(String[] args) throws Exception {
// 业务代码只需要用这个类,不用管后面的集群
SmartMQProducer producer = new SmartMQProducer();
producer.sendMessage("order-topic", "订单ID:12345");
}
}
这个SDK的核心逻辑是:自己每隔1秒检查当前用的集群是不是正常,如果不正常,就自动切换到备集群,对业务代码完全透明——业务只需要调用sendMessage方法,不用管哪个集群出问题了。
四、怎么规避脑裂和数据丢失?
刚才的方案看起来没问题,但有两个致命的坑:脑裂和数据丢失,必须提前规避。
4.1 什么是脑裂?怎么避免?
脑裂就是两个集群都觉得自己是主集群,同时接受业务的写请求,结果数据被写到两个集群里,两边的数据不一样,最后业务乱套。
举个例子:主集群和备集群之间的网络断了,主集群没挂,但备集群收不到主集群的心跳,以为主集群挂了,就把VIP抢过去了;同时主集群还在正常工作,结果两个集群都能接受写请求,数据两边各存一部分,最后完全对不上。
怎么避免脑裂?核心是让两个集群不能同时成为主集群,我们可以加两个限制:
- 加一个第三方的仲裁节点:比如用Redis或者ZooKeeper,只有拿到仲裁节点锁的集群,才能接受写请求。
- 备集群平时是只读的:只有切换过去之后,才变成可写的。
给刚才的Keepalived加个仲裁逻辑,用Redis做仲裁节点:
# 主集群的Keepalived配置,加个检查脚本,只有拿到Redis锁才当主
vrrp_script check_redis_lock {
script "/usr/local/bin/check_redis_lock.sh"
interval 2
weight 20
}
vrrp_instance VI_1 {
# ... 原来的配置不变
track_script {
check_redis_lock
}
}
# 检查Redis锁的脚本check_redis_lock.sh
#!/bin/bash
# 连接Redis,尝试拿锁,锁的过期时间设为5秒
LOCK_KEY="rocketmq-master-lock"
LOCK_VALUE="cluster-a"
# 尝试拿锁,只有拿到锁才返回0(正常)
redis-cli -h redis-server -p 6379 set $LOCK_KEY $LOCK_VALUE ex 5 nx
if [ $? -eq 0 ]; then
exit 0
else
exit 1
fi
备集群的检查脚本逻辑一样,只是LOCK_VALUE设为cluster-b。这样一来,只有拿到Redis锁的集群,才能当主集群,不会出现两个集群同时当主的情况,彻底避免脑裂。
4.2 什么是数据丢失?怎么避免?
数据丢失有两种情况:一是主集群挂了,还有一部分消息没同步到备集群,切换过去之后,这部分消息丢了;二是切换的时候,有消息正在发,结果发了一半断了,消息没存下来。
怎么避免?核心是保证数据同步的一致性,我们可以加两个机制:
- 主从同步用同步刷盘:主集群的Broker必须配置成同步刷盘,也就是消息写到主Broker之后,必须同步到从Broker,才给业务返回成功。
- 切换的时候做数据对齐:切换之前,先检查主集群的未同步消息,把这些消息同步到备集群,再切换。
给RocketMQ的Broker加同步刷盘的配置(在broker.conf里):
# 同步刷盘配置,消息必须同步到从Broker才返回成功
brokerRole=SYNC_MASTER
flushDiskType=SYNC_FLUSH
这个配置的意思是:业务发消息到主Broker,主Broker必须把消息同步到所有从Broker,才告诉业务“消息发成功了”。这样一来,主集群挂了,所有已经发成功的消息,都已经同步到从Broker,不会丢。
另外,切换的时候,我们可以加一个数据对齐的步骤:比如当主集群挂了,备集群抢了VIP之后,先把主集群的未同步消息拉过来,再接受业务的写请求。这个逻辑可以加在SDK的切换逻辑里,比如刚才的checkClusterHealth方法里,切换之前先调用一个alignData方法,把主集群的消息同步过来。
五、方案的优缺点和适用场景
5.1 优点
- 切换速度快:VIP切换不到3秒,SDK切换不到1秒,整个故障转移过程能控制在5秒以内,用户几乎感觉不到。
- 对业务透明:业务代码只需要写一个VIP,不用改任何配置,不用重启,运维成本低。
- 可靠性高:彻底避免脑裂和数据丢失,适合对数据一致性要求高的业务。
5.2 缺点
- 运维成本:需要维护两个一模一样的集群,还要维护Keepalived、Redis仲裁节点,运维的工作量比单集群大。
- 成本高:两个集群的硬件成本是单集群的两倍,适合预算充足的公司。
- 切换的时候有短暂的停顿:虽然时间很短,但还是有几毫秒的停顿,不适合对延迟要求极高的业务(比如金融交易的毫秒级延迟要求)。
5.3 适用场景
这个方案适合对数据一致性要求高、对可用性要求高的业务,比如订单系统、支付系统、生鲜配送系统、电商系统等;不适合对延迟要求极高、预算有限的业务。
六、方案落地的注意事项
- 两个集群的配置必须完全一致:包括NameServer的数量、Broker的数量、主题的配置、队列数、权限等,不然切换过去之后,业务会出问题。
- 仲裁节点必须可靠:比如用Redis做仲裁节点,Redis不能挂,不然整个容灾机制就失效了,所以Redis最好也做集群,保证高可用。
- 备集群平时是只读的:只有切换过去之后才变成可写的,不然两个集群都能写,还是会出脑裂的问题。
- 要做故障演练:平时要定期模拟集群故障,测试切换是不是正常,数据是不是不丢,不能等到真的出问题才发现容灾机制没用。
七、总结
RocketMQ多集群容灾的核心,是用VIP切换解决集群IP的问题,用SDK智能感知解决业务不用改配置的问题,同时加仲裁节点避免脑裂,加同步刷盘避免数据丢失。这个方案的本质,是让故障转移对业务透明,同时保证数据的一致性和可靠性。
最后要提醒大家,容灾不是一劳永逸的,要定期检查、定期演练,不然真的出问题的时候,容灾机制反而会变成新的故障点。
评论
围绕“构建RocketMQ多集群容灾架构时,通过VIP切换与SDK智能感知实现快速故障转移,需要规避脑裂与数据丢失风险的完整设计方案”参与讨论