前言
在现代微服务架构演进的过程中,系统往往会面临从“单机架构”向“高可用集群架构”升级的痛点。前阵子一位同行私信我,说他接手了一套核心业务系统,Redis、RabbitMQ、Elasticsearch 均部署在单节点上。随着业务量突增,系统面临极高的高可用与性能风险,迫切需要升级到集群模式。
但他面临的极大挑战在于:业务极其敏感,不能挂停机维护公告,必须在迁移全程保障业务零中断(7×24小时可用)。
更具体的约束条件是:
面对“不能停机、不能丢数据、不能重复”的苛刻要求,传统“选择深夜停机 -> 导出数据 -> 批量导入 -> 重新上线”的方案直接失效。
本文基于真实的生产实战经验,整理了一套针对 Redis、RabbitMQ、Elasticsearch 三大中间件从单节点平滑迁移至集群的无缝方案。全文摒弃对第三方插件的过度依赖,重点采用“纯应用层双写 + 消费端强幂等”的彻底自控逻辑,涵盖方案选型、配置代码、避坑细节、运维命令与回退预案,希望能为有类似需求的同行提供一份可复用的实战参考。
⚠️ 重要提示:本文方案为通用生产技术参考,所有代码与脚本在上线前务必先在测试环境完整演练,验证无误后再谨慎操作。因未充分测试或直接在生产环境误操作导致的任何业务风险与损失,本文作者概不负责。
一、 背景与核心挑战
1.1 典型单节点架构现状
许多中小型系统在初期演进中容易留下“单点隐患”:
| Redis | Single Node (无主从) | 宕机导致缓存全失效,极易引发数据库击穿与雪崩 |
| RabbitMQ | Single Broker (单节点) | 宕机导致异步消息链路断裂,存在单点丢失风险 |
| Elasticsearch | Single Node (单节点) | 无副本机制,节点故障导致检索服务瘫痪,存在数据损坏风险 |
1.2 迁移面临的三大难题
二、 迁移方案总体架构设计
在设计迁移策略时,我们需要针对不同中间件的数据特性,精确制定是否需要应用层双写:
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 的快照数据复制,存在两个致命缺陷:
因此,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"
}
}'
第五步:切换读流量与关闭双写
五、 RabbitMQ:单节点迁移至仲裁队列集群(应用双写:生产单写新集群 + 消费双监听 + 幂等防重)
5.1 为什么选择纯应用双写方案?
虽然 RabbitMQ 官方提供了 Federation(联邦)插件,但在许多生产环境中:
说明:本文的“双写”特指应用同时连接新旧两个 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 迁移实施标准步骤
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. 旧节点恢复独立对外服务
动态回退标准动作:
七、 生产迁移标准化 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 队列集群。
这套方案已经在高并发线上场景落地验证。希望这份指南能为你后续的架构演进与集群改造提供清晰、安全的落地路径!




