欢迎光临
我们一直在努力

业务不停机、数据零丢失:Redis / RabbitMQ / Elasticsearch 三大核心中间件升级改造集群无缝平滑迁移

前言

在现代微服务架构演进的过程中,系统往往会面临从“单机架构”向“高可用集群架构”升级的痛点。前阵子一位同行私信我,说他接手了一套核心业务系统,Redis、RabbitMQ、Elasticsearch 均部署在单节点上。随着业务量突增,系统面临极高的高可用与性能风险,迫切需要升级到集群模式。

但他面临的极大挑战在于:业务极其敏感,不能挂停机维护公告,必须在迁移全程保障业务零中断(7×24小时可用)。

更具体的约束条件是:

  • 旧节点上积压的待消费消息一条都不能丢失,且必须在新集群中平滑被消费;
  • 切换切流后,新产生的消息只能写入新集群,旧节点不再接收新数据;
  • 数据绝不能重复处理,必须做到严格的幂等保障。
  • 面对“不能停机、不能丢数据、不能重复”的苛刻要求,传统“选择深夜停机 -> 导出数据 -> 批量导入 -> 重新上线”的方案直接失效。

    本文基于真实的生产实战经验,整理了一套针对 Redis、RabbitMQ、Elasticsearch 三大中间件从单节点平滑迁移至集群的无缝方案。全文摒弃对第三方插件的过度依赖,重点采用“纯应用层双写 + 消费端强幂等”的彻底自控逻辑,涵盖方案选型、配置代码、避坑细节、运维命令与回退预案,希望能为有类似需求的同行提供一份可复用的实战参考。

    ⚠️ 重要提示:本文方案为通用生产技术参考,所有代码与脚本在上线前务必先在测试环境完整演练,验证无误后再谨慎操作。因未充分测试或直接在生产环境误操作导致的任何业务风险与损失,本文作者概不负责。


    一、 背景与核心挑战

    1.1 典型单节点架构现状

    许多中小型系统在初期演进中容易留下“单点隐患”:

    中间件当前单节点架构潜在瓶颈与风险
    Redis Single Node (无主从) 宕机导致缓存全失效,极易引发数据库击穿与雪崩
    RabbitMQ Single Broker (单节点) 宕机导致异步消息链路断裂,存在单点丢失风险
    Elasticsearch Single Node (单节点) 无副本机制,节点故障导致检索服务瘫痪,存在数据损坏风险

    1.2 迁移面临的三大难题

  • 业务零中断:迁移过程中,所有对外 API 与内部 RPC 请求不能出现拒绝服务或异常超时。
  • 消息零丢失 + 零重复:旧节点积压的待消费消息必须一条不漏地消化,且同一笔业务绝不能重复处理。
  • 新数据单向流转:切换完成后,新流量只能写入新集群,旧节点变为只读/只出不进状态。

  • 二、 迁移方案总体架构设计

    在设计迁移策略时,我们需要针对不同中间件的数据特性,精确制定是否需要应用层双写:

    2.1 三大中间件迁移策略矩阵表

    中间件是否需要应用双写核心迁移机制原因与保障逻辑
    Redis 不需要 RedisShake 增量同步 RedisShake 支持实时 AOF 增量同步,工具本身已做到物理级实时追平,应用层双写反易引发分布式锁与 TTL 冲突。
    Elasticsearch 强烈需要 _reindex 快照 + 应用双写 _reindex 无法感知物理删除(Delete)与并发更新,必须靠应用双写(或 Binlog 监听)实时覆盖新变更。
    RabbitMQ 强烈需要 应用双写(生产单写新集群 + 消费双监听)+ 去重表 不依赖 Federation 插件。生产者直接切为只写新集群,消费者同时监听新旧集群并靠 DB 唯一索引拦截重复,旧队列只出不进自然归零后下线旧节点。

    说明:

    • Redis 和 RabbitMQ 的迁移目标均为高可用集群架构(Redis 三主三从、RabbitMQ 三节点仲裁队列),因为这两个组件一旦宕机将直接影响业务数据流转。
    • Elasticsearch 的迁移目标可以根据自身架构需求选择单节点或集群,迁移方法完全相同,具体取决于你是否需要 ES 高可用。

    三、 Redis:单节点迁移至 Cluster 三主三从集群

    3.1 同步工具:RedisShake

    推荐使用阿里云开源的 RedisShake 3.x,它支持将 Redis 单节点数据以“全量 + 增量”的形式平滑同步至 Cluster 集群。

    3.2 迁移前开发侧必须排查的事项

    Redis Cluster 与单节点有本质区别:数据分片、客户端路由、跨 key 操作限制。切流前,开发同事必须完成以下排查与改造,否则切换后可能出现运行时错误。

    1. 客户端类型是否支持集群
    • 若使用 Jedis,需切换为 JedisCluster;
    • 若使用 Lettuce(Spring Boot 2.x 默认),天然支持集群,只需正确配置;
    • 若使用 Redisson,需使用集群模式 RedissonClient。

    不能继续使用单节点连接方式。

    2. 是否存在跨 Slot 的多 Key 操作

    集群模式下,mget、mset、del 多个 key、自定义 Lua 脚本、事务等跨多个 key 的操作会抛出:

    CROSSSLOT Keys in request don't hash to the same slot

    解决办法:

    • 使用 Hash Tag 强制相关 key 落在同一个 slot,例如将 user:100:profile 改为 {user:100}:profile,user:100:avatar 改为 {user:100}:avatar;
    • 或者拆分为多个单 key 操作(注意原子性可能受影响)。
    3. Pipeline 和事务
    • Pipeline 仍可用,但需客户端支持自动路由;
    • 事务(MULTI/EXEC)只能作用于同一个 slot 内的 key,跨 slot 会失败。
    4. 客户端配置
    • 配置所有 master 节点地址(或至少配置一个,客户端自动发现拓扑);
    • 设置 max-redirects(如 spring.redis.cluster.max-redirects=3);
    • 确保密码、超时等配置正确。

    建议在测试环境模拟集群连接后跑一遍核心业务回归。

    3.3 实战操作步骤

    第一步:部署新 Redis Cluster 集群

    在目标服务器部署三主三从集群,并确保版本与旧节点一致或向后兼容。

    # 创建 3 主 3 从集群(示例)
    redis-cli –cluster create \\
    node1:6379 node2:6379 node3:6379 \\
    node1:6380 node2:6380 node3:6380 \\
    –cluster-replicas 1 -a "your_password"

    第二步:配置并启动 RedisShake

    解压并修改配置文件 redis-shake.toml:

    [source]
    type = "standalone"
    address = "192.168.1.10:6379" # 旧 Redis 单节点 IP
    password = "old_password"

    [target]
    type = "cluster"
    address = "192.168.1.20:6379" # 新集群任意 Master 节点 IP
    password = "new_password"

    [sync_reader]
    cluster = false

    [sync_writer]
    cluster = true

    启动同步进程:

    ./redis-shake sync –config redis-shake.toml

    第三步:数据一致性校验

    利用 redis-full-check 工具校验新旧节点 Key 数量及内容一致性:

    redis-full-check -s 192.168.1.10:6379 -t 192.168.1.20:6379 -a "your_password" –comparemode=1

    第四步:配置中心动态切流

    同步状态达到实时增量后,在 Nacos 中更新 Redis 连接配置,业务应用感知配置变更后平滑连接新集群:

    # 原旧节点配置
    # spring.redis.host=192.168.1.10
    # spring.redis.port=6379

    # 新 Cluster 配置
    spring.redis.cluster.nodes=192.168.1.20:6379,192.168.1.21:6379,192.168.1.22:6379
    spring.redis.cluster.max-redirects=3
    spring.redis.password=new_password


    四、 Elasticsearch:单节点迁移至集群或单节点(Reindex + 应用双写)

    4.1 为什么 ES 必须配合应用双写?

    Elasticsearch 内置的 _reindex API 本质上是基于 Scroll 的快照数据复制,存在两个致命缺陷:

  • 无法同步物理删除(Delete):若在 Reindex 执行期间旧 ES 删除了某条文档,Reindex 无法感知此操作,导致新 ES 留存废弃脏数据。
  • 并发更新覆盖:海量数据 Reindex 耗时可能数小时,这期间产生的业务更新极易被 Reindex 批处理覆盖为旧版本。
  • 因此,ES 必须使用“应用双写 + 存量 _reindex(Ignore 冲突)”配合迁移。

    迁移目标说明:本方案适用于旧 ES 单节点迁移到新 ES 单节点或集群,操作方法完全一致。如果你的业务对检索服务可用性要求较高(例如核心交易搜索),建议目标采用三节点集群;如果仅作为日志或辅助检索,单节点即可。关键点在于应用双写保证增量数据一致性。


    4.2 实战操作步骤

    第一步:导出并手动重建 Mapping / Settings

    提前在目标新 ES(单节点或集群)上手动创建索引及 Mapping,切忌依赖 ES 自动推断(防止 keyword 被推断为 text):

    #!/bin/bash
    OLD_ES="http://admin:password@192.168.1.10:9200"
    NEW_ES="http://admin:password@192.168.1.20:9200"

    for INDEX in $(curl -s "${OLD_ES}/_cat/indices/payment-*?h=index"); do
    # 1. 提取旧 Mapping 与 Settings
    curl -s "${OLD_ES}/${INDEX}/_mapping" | jq ".\\"${INDEX}\\".mappings" > "/tmp/${INDEX}_mapping.json"
    curl -s "${OLD_ES}/${INDEX}/_settings" | jq ".\\"${INDEX}\\".settings | {index: {analysis: .index.analysis}}" > "/tmp/${INDEX}_settings.json"

    # 2. 在新集群上提前创建对应索引
    curl -s -XPUT "${NEW_ES}/${INDEX}" \\
    -H 'Content-Type: application/json' \\
    -d "{\\"mappings\\": $(cat /tmp/${INDEX}_mapping.json), \\"settings\\": $(cat /tmp/${INDEX}_settings.json)}"
    done

    第二步:开启 ES 应用层双写(关键)

    在业务代码中对文档的增、删、改逻辑引入双写服务:

    @Service
    @Slf4j
    public class EsDualWriteService {

    @Autowired
    @Qualifier("oldEsClient")
    private RestHighLevelClient oldEsClient;

    @Autowired
    @Qualifier("newEsClient")
    private RestHighLevelClient newEsClient;

    @Value("${config.es.write-mode:OLD_ONLY}") // OLD_ONLY | DUAL | NEW_ONLY
    private String esWriteMode;

    public void saveDocument(String index, String id, Map<String, Object> data) {
    IndexRequest request = new IndexRequest(index).id(id).source(data);

    // 写旧 ES
    if ("OLD_ONLY".equals(esWriteMode) || "DUAL".equals(esWriteMode)) {
    try { oldEsClient.index(request, RequestOptions.DEFAULT); } catch (Exception e) { log.error("旧 ES 写入失败", e); }
    }

    // 双写新 ES(单节点或集群)
    if ("NEW_ONLY".equals(esWriteMode) || "DUAL".equals(esWriteMode)) {
    try { newEsClient.index(request, RequestOptions.DEFAULT); } catch (Exception e) { log.error("新 ES 双写失败", e); }
    }
    }

    // 删除动作双写:彻底解决 _reindex 无法同步删除的问题
    public void deleteDocument(String index, String id) {
    DeleteRequest request = new DeleteRequest(index, id);
    if ("OLD_ONLY".equals(esWriteMode) || "DUAL".equals(esWriteMode)) {
    try { oldEsClient.delete(request, RequestOptions.DEFAULT); } catch (Exception ignored) {}
    }
    if ("NEW_ONLY".equals(esWriteMode) || "DUAL".equals(esWriteMode)) {
    try { newEsClient.delete(request, RequestOptions.DEFAULT); } catch (Exception ignored) {}
    }
    }
    }

    第三步:配置远程 Reindex 白名单(重要)

    由于 _reindex 需要从旧 ES 远程读取数据,必须在目标新 ES 节点的 elasticsearch.yml 中添加白名单,否则会报错:

    # 在目标新 ES 节点配置,允许多个远程源用逗号分隔
    reindex.remote.whitelist: "192.168.1.10:9200"

    修改后需要滚动重启目标新 ES 节点(若为集群,逐个重启以保持可用性)。

    第四步:执行全量快照 Reindex(op_type = create)

    在应用双写开启后,启动 Reindex 搬运历史存量数据。关键设置 "op_type": "create"(或忽略冲突):这样如果某条数据已经被应用双写实时更新到了新 ES 中,Reindex 就不会用历史旧数据去覆写它!

    curl -XPOST 'http://192.168.1.20:9200/_reindex?wait_for_completion=false&slices=auto&requests_per_second=500' \\
    -H 'Content-Type: application/json' \\
    -d '{
    "conflicts": "proceed",
    "source": {
    "remote": {
    "host": "http://192.168.1.10:9200",
    "username": "admin",
    "password": "password"
    },
    "index": "payment-*",
    "size": 1000
    },
    "dest": {
    "index": "payment-*",
    "op_type": "create"
    }
    }'

    第五步:切换读流量与关闭双写
  • 确认 Reindex 任务完成,且应用双写无报错;
  • 切换应用查询 Client 连接指向新 ES(单节点或集群);
  • 将 Nacos 中 es.write-mode 改为 NEW_ONLY,停止向旧 ES 写入。

  • 五、 RabbitMQ:单节点迁移至仲裁队列集群(应用双写:生产单写新集群 + 消费双监听 + 幂等防重)

    5.1 为什么选择纯应用双写方案?

    虽然 RabbitMQ 官方提供了 Federation(联邦)插件,但在许多生产环境中:

  • 运维严禁在生产 MQ 实例上安装/配置额外的第三方插件;
  • 插件搬运消息与业务消费者存在抢占竞争,极易导致重复消费;
  • 应用代码自控是唯一能 100% 确保“老积压消息一条不丢、新消息只落新集群、业务强幂等去重”的可靠路径。
  • 说明:本文的“双写”特指应用同时连接新旧两个 MQ 集群,但生产者的行为是单写新集群,只有消费者是双监听。这样旧 MQ 只出不进,历史消息自然清空,不会出现“旧队列永远无法归零”的问题。


    5.2 架构演进全流程

    【阶段 1:只写只读旧节点】
    生产者 ──► 旧 MQ (积压老消息) ──► 消费者 (连旧 MQ)

    【阶段 2:生产单写新集群 + 消费双监听】
    生产者 ──► (写开关: NEW_ONLY) ──► 新 MQ ──► 消费者 (监听新 MQ) ──┐
    旧 MQ ──► 消费者 (监听旧 MQ) ──┼─► [去重表强拦截]

    【阶段 3:旧堆积归零 + 下线旧节点】 │
    旧 MQ 堆积归零 ──► 停止旧监听器 ──► 下线旧 MQ │

    (旧 MQ 只出不进,自然清空) │


    5.3 实战完整代码实现

    1. 生产者:动态写控制组件(支持单写新集群)

    通过 Nacos 配置中心下发 write-mode 属性,实现对写行为的动态切换:

    @Component
    @RefreshScope // 支持 Nacos / Apollo 动态刷新配置
    @Slf4j
    public class RabbitMqPublisher {

    // OLD_ONLY (仅写旧) | NEW_ONLY (仅写新)
    // 迁移切换时直接从 OLD_ONLY 切到 NEW_ONLY,无需 DUAL
    @Value("${config.rabbitmq.write-mode:OLD_ONLY}")
    private String writeMode;

    @Autowired
    @Qualifier("oldRabbitTemplate")
    private RabbitTemplate oldRabbitTemplate;

    @Autowired
    @Qualifier("newRabbitTemplate")
    private RabbitTemplate newRabbitTemplate;

    public void send(String exchange, String routingKey, Object message) {
    // 根据配置决定写哪个集群
    if ("OLD_ONLY".equals(writeMode)) {
    try {
    oldRabbitTemplate.convertAndSend(exchange, routingKey, message);
    } catch (Exception e) {
    log.error("写入旧 RabbitMQ 失败", e);
    // 根据业务决定是否抛出异常
    }
    } else if ("NEW_ONLY".equals(writeMode)) {
    try {
    newRabbitTemplate.convertAndSend(exchange, routingKey, message);
    } catch (Exception e) {
    log.error("写入新 RabbitMQ 失败,需投递死信或异步补偿", e);
    // 强烈建议接入本地消息表或异步补偿任务,保证最终一致
    }
    }
    }
    }

    2. 消费者端:双集群监听 + 数据库幂等去重表

    迁移期间,新旧集群都会存在消息。利用数据库联合主键/唯一索引构建去重表,拦截重复消费:

    — 消费幂等去重表
    CREATE TABLE `sys_msg_dedup` (
    `tx_id` varchar(64) NOT NULL COMMENT '业务流水/消息唯一ID',
    `consumer_group` varchar(64) NOT NULL COMMENT '消费组标识',
    `created_at` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (`tx_id`, `consumer_group`)
    ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

    @Component
    @Slf4j
    public class PaymentMessageConsumer {

    @Autowired
    private MsgDedupMapper dedupMapper;

    @Autowired
    private PaymentService paymentService;

    // 监听旧节点队列
    @RabbitListener(queues = "${config.old.queue.name}", containerFactory = "oldContainerFactory")
    public void handleOldClusterMessage(PaymentMessage msg, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
    processWithDedup(msg, channel, tag, "OLD_MQ");
    }

    // 监听新集群 Quorum 队列
    @RabbitListener(queues = "${config.new.queue.name}", containerFactory = "newContainerFactory")
    public void handleNewClusterMessage(PaymentMessage msg, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
    processWithDedup(msg, channel, tag, "NEW_MQ");
    }

    private void processWithDedup(PaymentMessage msg, Channel channel, long tag, String source) throws IOException {
    String txId = msg.getTxId(); // 获取消息绑定的业务唯一 ID
    try {
    // 1. 尝试插入去重表 (txid + consumer_group 为联合主键)
    dedupMapper.insert(new MsgDedupRecord(txId, "payment_service_group"));

    // 2. 插入成功,说明未被消费过,执行核心业务
    paymentService.doPayment(msg);

    // 3. 业务执行成功,手动 ACK
    channel.basicAck(tag, false);
    log.info("消息消费成功 [来源: {}], txId: {}", source, txId);

    } catch (DuplicateKeyException e) {
    // 4. 触发唯一索引冲突:说明另一套集群传来的相同消息已被处理过!直接 ACK 丢弃
    log.warn("检测到重复消息 [来源: {}],已通过幂等去重表拦截,txId: {}", source, txId);
    channel.basicAck(tag, false);

    } catch (Exception e) {
    log.error("业务处理失败 [来源: {}], txId: {}", source, txId, e);
    // 5. 业务异常时,务必先删除去重表记录,再拒绝消息重新入队,避免重试时被误判为重复消息
    try {
    dedupMapper.delete(txId, "payment_service_group");
    } catch (Exception deleteEx) {
    log.error("删除去重表记录失败,可能导致消息无法重试, txId: {}", txId, deleteEx);
    // 此处可考虑抛出异常阻断 ACK,或投递死信队列人工介入
    }
    // 6. 拒绝并重新入队,等待下次重试
    channel.basicNack(tag, false, true);
    }
    }
    }

    核心改进说明:业务处理失败时,去重表记录会被删除,这样消息重新入队后再次消费时,插入去重表不会冲突,业务得以继续执行。若删除失败,需接入死信队列或告警,防止消息永久丢失。


    5.4 迁移实施标准步骤

  • 部署新 Quorum 集群:搭建三节点仲裁队列集群,提取旧节点 Queue 定义,注入 x-queue-type: quorum 后在新集群提前创建 Exchange、Queue 与 Binding。
  • 上线双监听消费者:发布代码,让消费者同时监听新旧两个 MQ 的队列。确保去重表逻辑正确生效。
  • 切换生产者写模式为 NEW_ONLY:通过 Nacos 将 write-mode 从 OLD_ONLY 改为 NEW_ONLY。此时新消息只进新 MQ,旧 MQ 不再接收新消息。
  • 监控旧队列归零:使用命令持续监控旧节点:
  • watch -n 3 'rabbitmqctl list_queues name messages messages_unacknowledged'

    旧 MQ 只出不进,messages 与 messages_unacknowledged 会持续下降,最终完全降为 0。
    5. 安全下线:确认旧 MQ 队列归零后,停止旧集群监听器,关闭并下线旧单节点 RabbitMQ。


    六、 配置中心动态切流与秒级回退预案

    任何生产级别的迁移方案,都必须设计“秒级回退预案”。

    ┌─────────────────────────┐
    │ Nacos / Apollo 配置中心 │
    └────────────┬────────────┘

    ┌───────────────┴───────────────┐
    ▼ ▼
    【正常切流路径】 【秒级回退路径】
    1. write-mode = NEW_ONLY 1. write-mode = OLD_ONLY
    2. 部署双监听消费者 2. 切换各连接配置回旧节点 IP
    3. 旧 MQ 只出不进,自然清空 3. 关闭新集群入口流量
    4. 确认旧队列归零后下线 4. 旧节点恢复独立对外服务

    动态回退标准动作:

  • 下发回退指令:将 write-mode 统一改回 OLD_ONLY;
  • 切回连接配置:将 Redis、RabbitMQ、ES 连接地址一键切回旧单节点 IP;
  • 数据完好无损:在整个迁移过程中,旧节点数据全程保留,并未执行物理删除,系统可以在秒级恢复原状。

  • 七、 生产迁移标准化 Checklist

    阶段核心任务状态确认
    准备阶段 新集群环境搭建完毕(Redis Cluster / MQ Quorum / ES 单节点或集群) [ ]
    完成代码改造:Redis 客户端集群支持;MQ/ES 动态双写与 sys_msg_dedup 去重表 [ ]
    ES 新索引 Mapping 手动创建完成;RabbitMQ 提前创建 Quorum 队列 [ ]
    开发侧排查 Redis 跨 slot 操作,必要时加入 Hash Tag [ ]
    同步阶段 启动 RedisShake,确认增量同步延时降至 0ms [ ]
    开启 ES 应用双写(write-mode = DUAL),配置远程 Reindex 白名单并重启目标节点 [ ]
    执行 ES 全量 _reindex(设置 op_type=create 避免覆盖双写) [ ]
    切流阶段 部署 RabbitMQ 双监听消费者,验证数据库去重表拦截有效 [ ]
    将 RabbitMQ 生产者写模式切换为 NEW_ONLY,停止向旧 MQ 发送消息 [ ]
    监控旧 MQ 堆积量持续下降直至归零 [ ]
    切换 Redis 与 ES 读连接至新集群/新节点 [ ]
    收尾阶段 新集群 CPU/内存/IO/响应延时(P99)等指标表现平稳 [ ]
    停止 RedisShake 进程,清理旧节点监听器 [ ]
    保留旧节点冷备数据 7 天后,安全下线旧单节点 [ ]

    八、 总结

    线上系统的中间件平滑升级,看似是工具的使用,本质上是对流量控制与数据状态变化的精细化掌控。

    针对本文的三大核心中间件,我们总结出这套极其扎实的实战策略:

    • Redis:依靠 RedisShake 物理级增量同步,应用层需适配集群客户端与跨 slot 操作,目标是三主三从 Cluster;
    • Elasticsearch:采用应用双写(处理删改) + 存量 _reindex,彻底抹平并发更新冲突,目标可以是单节点或集群,视高可用需求而定;
    • RabbitMQ:采用应用双写(生产单写新集群 + 消费双监听)+ DB 去重拦截,摆脱插件依赖,绝对保障老消息不丢、新消息只发新集群,目标是三节点 Quorum 队列集群。

    这套方案已经在高并发线上场景落地验证。希望这份指南能为你后续的架构演进与集群改造提供清晰、安全的落地路径!
    在这里插入图片描述

    赞(0)
    未经允许不得转载:171主机测评 » 业务不停机、数据零丢失:Redis / RabbitMQ / Elasticsearch 三大核心中间件升级改造集群无缝平滑迁移
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址