一、为什么要做RocketMQ多集群容灾?

先给大家说个真实的坑:去年我帮一家做生鲜配送的公司排查故障,他们的订单系统全靠RocketMQ存订单消息,结果某天机房空调坏了,整个RocketMQ集群直接挂了。更坑的是,他们当时只有这一个集群,挂了之后订单全乱了——有的用户付了钱没收到确认,有的骑手拿不到配送单,整整折腾了4个小时才恢复,赔了客户几十万。

这就是单集群的致命问题:只要集群出问题,整个依赖它的业务都得停。多集群容灾的核心,就是不让一个集群的故障波及所有业务,而且故障发生时,能快速切换到正常的集群,尽量不影响用户体验。

二、核心思路:VIP切换+SDK智能感知

很多人做容灾的第一反应是“集群挂了就切IP”,但直接改IP有两个大问题:一是业务代码要手动改配置重启,慢得要死;二是如果集群没完全挂,只是部分节点出问题,乱切换反而会出乱子。

我们的方案是把两个核心技术结合起来:

  1. VIP切换:给两个RocketMQ集群绑同一个虚拟IP(VIP),平时只有主集群用这个IP,主集群出问题时,自动把VIP切到备集群;
  2. 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抢过去了;同时主集群还在正常工作,结果两个集群都能接受写请求,数据两边各存一部分,最后完全对不上。

怎么避免脑裂?核心是让两个集群不能同时成为主集群,我们可以加两个限制:

  1. 加一个第三方的仲裁节点:比如用Redis或者ZooKeeper,只有拿到仲裁节点锁的集群,才能接受写请求。
  2. 备集群平时是只读的:只有切换过去之后,才变成可写的。

给刚才的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 什么是数据丢失?怎么避免?

数据丢失有两种情况:一是主集群挂了,还有一部分消息没同步到备集群,切换过去之后,这部分消息丢了;二是切换的时候,有消息正在发,结果发了一半断了,消息没存下来。

怎么避免?核心是保证数据同步的一致性,我们可以加两个机制:

  1. 主从同步用同步刷盘:主集群的Broker必须配置成同步刷盘,也就是消息写到主Broker之后,必须同步到从Broker,才给业务返回成功。
  2. 切换的时候做数据对齐:切换之前,先检查主集群的未同步消息,把这些消息同步到备集群,再切换。

给RocketMQ的Broker加同步刷盘的配置(在broker.conf里):

# 同步刷盘配置,消息必须同步到从Broker才返回成功
brokerRole=SYNC_MASTER
flushDiskType=SYNC_FLUSH

这个配置的意思是:业务发消息到主Broker,主Broker必须把消息同步到所有从Broker,才告诉业务“消息发成功了”。这样一来,主集群挂了,所有已经发成功的消息,都已经同步到从Broker,不会丢。

另外,切换的时候,我们可以加一个数据对齐的步骤:比如当主集群挂了,备集群抢了VIP之后,先把主集群的未同步消息拉过来,再接受业务的写请求。这个逻辑可以加在SDK的切换逻辑里,比如刚才的checkClusterHealth方法里,切换之前先调用一个alignData方法,把主集群的消息同步过来。

五、方案的优缺点和适用场景

5.1 优点

  1. 切换速度快:VIP切换不到3秒,SDK切换不到1秒,整个故障转移过程能控制在5秒以内,用户几乎感觉不到。
  2. 对业务透明:业务代码只需要写一个VIP,不用改任何配置,不用重启,运维成本低。
  3. 可靠性高:彻底避免脑裂和数据丢失,适合对数据一致性要求高的业务。

5.2 缺点

  1. 运维成本:需要维护两个一模一样的集群,还要维护Keepalived、Redis仲裁节点,运维的工作量比单集群大。
  2. 成本高:两个集群的硬件成本是单集群的两倍,适合预算充足的公司。
  3. 切换的时候有短暂的停顿:虽然时间很短,但还是有几毫秒的停顿,不适合对延迟要求极高的业务(比如金融交易的毫秒级延迟要求)。

5.3 适用场景

这个方案适合对数据一致性要求高、对可用性要求高的业务,比如订单系统、支付系统、生鲜配送系统、电商系统等;不适合对延迟要求极高、预算有限的业务。

六、方案落地的注意事项

  1. 两个集群的配置必须完全一致:包括NameServer的数量、Broker的数量、主题的配置、队列数、权限等,不然切换过去之后,业务会出问题。
  2. 仲裁节点必须可靠:比如用Redis做仲裁节点,Redis不能挂,不然整个容灾机制就失效了,所以Redis最好也做集群,保证高可用。
  3. 备集群平时是只读的:只有切换过去之后才变成可写的,不然两个集群都能写,还是会出脑裂的问题。
  4. 要做故障演练:平时要定期模拟集群故障,测试切换是不是正常,数据是不是不丢,不能等到真的出问题才发现容灾机制没用。

七、总结

RocketMQ多集群容灾的核心,是用VIP切换解决集群IP的问题,用SDK智能感知解决业务不用改配置的问题,同时加仲裁节点避免脑裂,加同步刷盘避免数据丢失。这个方案的本质,是让故障转移对业务透明,同时保证数据的一致性和可靠性。

最后要提醒大家,容灾不是一劳永逸的,要定期检查、定期演练,不然真的出问题的时候,容灾机制反而会变成新的故障点。