做集群迁移这事,最怕的不是数据量大,而是迁移过程中业务还在跑,存量数据和新写入的数据搅在一起,割接时又发现消费端衔接不上,最后变成一个没法收场的“半迁移”状态。我踩过这个坑之后,遇到跨集群搬迁的需求,第一反应就是整套MirrorMaker2的方案走一遍:它有现成的复制通道、消费组位点同步机制、心跳检测和自动发现topic的能力,虽然配置上有些细节容易踩雷,但整体比自研迁移工具、或者用旧版MirrorMaker肉搏靠谱太多。
这篇东西适合谁看?手上有Kafka集群要换机房、升版本、拆集群的运维和开发,还在纠结“双写 vs 迁移工具”怎么选的人,以及第一次接触MirrorMaker2、想直接照着配置落地的同学。我会把从环境准备、配置设计、启动验证,到消费者切换、问题排查的完整过程拆开写,所有参数都给到能直接抄作业的程度,并把我在实际操作中遇到的坑一并标注出来。
1. 迁移方案选型:为什么是MirrorMaker2而不是其他方案
1.1 先盘一下市面上的主流迁移思路
Kafka集群迁移,常见的有三条路。第一是业务双写,就是生产者同时往新旧集群发数据,跑一段时间后切读。这个方案逻辑最简单,但改动面很大,所有生产端代码都要动,而且双写期间消息顺序和幂等性都要额外保证,对业务团队来说是实打实的侵入性改造。第二是自研复制工具,自己写消费者拉旧集群数据、再生产到新集群,好处是灵活,坏处是offset管理、topic自动发现、分区映射这些全要自己处理,稍不留神就出偏差,时间成本也不低。第三就是用Kafka官方提供的MirrorMaker2(MM2),它本质上是Kafka Connect的一个连接器,基于配置就能完成跨集群复制,同时能同步消费组位点,这是它相比自研方案最大的优势——不需要业务配合,只在基础设施层操作,就能把整个集群的topic和数据完完整整搬到新集群去。
在真实的生产迁移场景里,我基本不会考虑双写自研的方案,除非公司有专门的中间件团队能长期维护一套迁移框架。大多数情况下,业务方给到的迁移窗口只有几个小时,这时候用MM2这种“开箱即用、可配置可监控”的方案,才是性价比最高的选择。另外有些人会拿旧版MirrorMaker(MM1)来做对比,这里必须强调一点:MM1只做数据镜像,不处理消费组位点迁移,而且多线程复制时分区顺序可能错乱,跨集群同步Topic配置和ACL更是无从谈起。MM2的诞生就是在补这些短板,所以新项目一律不要用MM1。
1.2 MirrorMaker2的核心机制决定了它适合迁移场景
MM2跑在Kafka Connect框架里,每个数据复制任务至少包含三个线程:AdminThread负责发现源集群的topic、消费组、配置变更并同步到目标集群;ConsumerThread从源集群拉取消息;ProducerThread把消息写入目标集群。这三个线程协同工作,让数据像是“水龙头对水龙头”一样流过去,而不是通过一个中间的存储层倒腾。
它真正厉害的地方在于内部设计了几个专用topic来支撑元数据同步。比如heartbeat topic,每个MM2节点会定时向目标集群发送心跳,用来判断复制链路是否健康;checkpoints topic记录源集群消费组的位点快照,目标集群侧可以据此对齐消费位置;offset syncs topic则保存源集群分区offset到目标集群分区offset的映射关系。这些topic在配置好MM2后会自动创建,不需要人工干预。理解了这个机制,你就明白为什么说MM2比其他自研方案适合做集群迁移了——它把迁移过程中最麻烦的“位点对齐”问题沉淀成了标准化的能力,不只是搬数据,是连着整个消费生态一起搬。
注意:MM2默认的命名规则是
{source.cluster.alias}->{target.cluster.alias}这样的双向或单向复制组合,内部topic名称也会带上alias前缀。如果新旧集群的别名取不好,后面排查问题时看topic列表会非常痛苦,这一点在后面配置章节会详细展开。
2. 迁移前的准备:版本、资源与拓扑设计
2.1 版本兼容性和运行环境要提前确认
开始动手之前,先把版本这个最大的变量定死。MM2不是一个独立软件,它是Kafka 2.4开始内置在Kafka发布包里的连接器,也就是说你下载Kafka发行版,里面就带MM2的jar包。但涉及跨集群迁移时,源集群和目标集群的Kafka版本可能不一样,这时候以谁的版本来跑Connect集群?我的经验是:Connect集群的版本尽量和目标集群保持一致,或者不低于源集群。因为MM2要同时连两个集群,以新版Kafka的客户端去兼容旧版broker通常问题不大,反之旧版客户端连新版broker就可能遇到协议不兼容的问题。
实操中我们需要准备一台独立的机器(或者一组机器)来跑Kafka Connect集群,不建议和业务broker混部。因为迁移期间数据复制流量很大,如果和业务broker混在一起,共享磁盘和带宽容易互相拖累;而且Connect重启、扩缩容也容易影响到broker进程。机器配置方面,CPU和内存主要取决于复制吞吐量,一般8核16G起步,磁盘几乎不需要——因为Connect是流式处理,不需要持久化数据。
2.2 命名规范:alias决定了整个迁移的“坐标系”
MM2里最容易被忽略但影响最深远的参数就是集群别名(alias)。每个集群都要给一个唯一别名,它会体现在复制通道名称、内部topic名称里。例如源集群叫oldCluster,目标集群叫newCluster,那么启动的复制任务就是oldCluster->newCluster,心跳topic是oldCluster.heartbeats,checkpoints topic是oldCluster.checkpoints.internal。这些名字一旦定下来,整个迁移过程中所有工具都会用到,所以别名选一个简短、语义明确的词最合适,比如src和dst,或者cluster-a、cluster-b。
我见过有团队把别名写成IP地址的,结果内部topic名字变成10-0-0-1.heartbeats,排查问题的时候一眼根本看不出来是心跳topic,特别膈应。规范的别名能让你在kafka-topics列表里一眼识别出哪些是MM2创建的内部topic,避免后面清理数据时误删。
2.3 需要提前梳理的清单:topic清单、消费组清单、分区策略
在写配置之前,列一个清单出来:
- 全量topic清单:用
kafka-topics --list拉出来,确认哪些topic需要迁移、哪些可以废弃、哪些包含敏感数据需要过滤。 - 消费组清单:用
kafka-consumer-groups --list拉出来,记录每个消费组的消费位点,尤其是Lag情况,这是迁移完成后做位点校验的基准。 - 分区和副本情况:记录每个topic的分区数、副本数、min.insync.replicas、cleanup.policy等关键配置,MM2其实会自动帮我们同步大部分topic配置,但保留原始记录总是更稳妥。
- 确认源集群是否开启了ACL和SSL:如果开启了,MM2连接源集群和目标集群时都要带上对应的安全认证参数,这一个项漏掉,会导致复制任务频繁报错。
这些信息整理成表格后,迁移结束的验收过程就有据可依。我不建议拿着topic列表现场核对,数据量大的时候容易看漏。
3. 核心实操:配置文件、启动和验证的完整步骤
3.1 第一步:搭建Kafka Connect运行环境
假设我们用Kafka 3.x版本自带的Connect,那么直接找到下载好的Kafka包,编辑config/connect-distributed.properties:
# Connect集群的组ID,注意不要和业务消费组重叠 group.id=mm2-cluster # Connect实例的id,每个节点唯一 client.id=mm2-connect-1 # 存储connector配置、offset、状态的topic,提前创建好或用auto.create config.storage.topic=mm2-configs offset.storage.topic=mm2-offsets status.storage.topic=mm2-status # 这几个内部topic的副本数,建议和集群保持一致 config.storage.replication.factor=3 offset.storage.replication.factor=3 status.storage.replication.factor=3 # 指定目标集群或可同时访问两个集群的bootstrap地址 bootstrap.servers=192.168.1.10:9092 # key/value converter建议用json,调试更方便 key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false value.converter.schemas.enable=false # REST接口,用来提交connector和管理任务 rest.port=8083 rest.advertised.host.name=192.168.1.20 rest.advertised.port=8083这里有几个非常重要的点。第一,config.storage.topic、offset.storage.topic、status.storage.topic这三个topic不要和MM2的复制topic混在一起,命名上加个统一前缀(如mm2-)最容易区分。第二,这些内部topic的分区数不宜太多,默认单分区或少量分区就行,因为Connect写入这些topic的频率不高。第三,分布式模式下,多个Connect节点会组成集群,connector任务会分布在不同节点上执行,所以节点数可以根据复制吞吐量水平扩展。如果只是临时迁移,跑单节点也完全够用。
启动命令很简单:
bin/connect-distributed.sh config/connect-distributed.properties启动后通过REST接口确认状态:
curl http://192.168.1.20:8083/出现{"version":"3.x.x","commit":"..."}就说明Connect起来了。
3.2 第二步:编写MirrorMaker2的connector配置文件
以JSON格式提交给Connect REST接口。这是我最常用的一份配置模板,几乎每个迁移项目都从它改出来的:
{ "name": "mm2-migrate-src-to-dst", "config": { "connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "tasks.max": "4", "source.cluster.alias": "src", "target.cluster.alias": "dst", "source.cluster.bootstrap.servers": "192.168.1.10:9092,192.168.1.11:9092", "target.cluster.bootstrap.servers": "192.168.2.10:9092,192.168.2.11:9092", "topics": ".*", "groups": ".*", "replication.factor": "3", "refresh.topics.interval.seconds": "300", "refresh.groups.interval.seconds": "300", "sync.topic.configs.enabled": "true", "sync.topic.acls.enabled": "false", "emit.heartbeats.interval.seconds": "5", "emit.checkpoints.interval.seconds": "5", "source.cluster.offset.syncs.topic.replication.factor": "3", "checkpoint.topic.replication.factor": "3", "heartbeats.topic.replication.factor": "3", "replication.policy.class": "org.apache.kafka.connect.mirror.DefaultReplicationPolicy" } }逐个说一下关键参数的选择理由:
tasks.max=4:这个值决定了数据复制任务拆成多少个task并行执行。MM2会按topic分区数自动做负载均衡,task越多并行度越高,但也不要无脑调高,因为每个task都会占用独立的连接和线程资源,建议先按topic分区总数除以100左右估算,后续观察consumer组的Lag再调整。topics=.*和groups=.*:用正则匹配所有topic和消费组。如果有特殊topic(比如业务内部的高频日志topic)不想要,可以写成topics=.*配合topics.exclude=internal_.*来过滤。refresh.topics.interval.seconds=300:每5分钟检查一次源集群有没有新增topic,发现后自动在目标集群创建并开始复制。这个值太小会频繁请求源集群元数据,太大会导致新增topic不能及时同步,5分钟是合理折中。refresh.groups.interval.seconds=300:同步消费组信息的时间间隔,同样影响位点同步的及时性。emit.checkpoints.interval.seconds=5:MM2每隔5秒把源集群消费组的offset快照写入checkpoints topic。这个值决定了迁移过程中,目标集群能多快感知到源集群消费位点的变化。时间设得太大,消费者切换时位点偏差就大。replication.factor=3:在目标集群创建topic时使用的副本数。这里很容易踩坑:如果你把复制因子设得比目标集群的最小ISR还低,那么同步过去后生产的可用性会出问题;如果你设得比目标集群broker数还高,topic创建直接失败。所以务必根据目标集群的节点数来定。sync.topic.configs.enabled=true:自动把源集群topic的配置同步到目标集群,比如retention.ms、cleanup.policy这些,这个功能很实用,但要注意,如果源集群有些topic配置比较特殊(比如无限保留、超大分区),同步过去后可能会影响目标集群的存储,所以有特殊topic建议在topics.exclude里排除掉,或者迁移后手动修正再放开同步。sync.topic.acls.enabled=false:如果源集群没开ACL,保持false;如果开了ACL,必须配true,否则复制过程中权限相关元数据不会同步,目标集群会丢授权信息。
提示:MirrorSourceConnector和MirrorCheckpointConnector是两回事。上面的配置用的是MirrorSourceConnector,它是数据复制的主力。MirrorCheckpointConnector是专门同步消费组位点的,在有些版本里需要单独再启动一个connector。我自己在用的Kafka 3.x版本里,MirrorSourceConnector已经默认包含消费组位点同步,所以只配一个就够。如果发现消费者切过去后位点不对,再检查是否需要单独加MirrorCheckpointConnector。
提交connector的命令:
curl -X POST http://192.168.1.20:8083/connectors \ -H "Content-Type: application/json" \ -d @mm2-config.json查看connector状态:
curl http://192.168.1.20:8083/connectors/mm2-migrate-src-to-dst/status正常的情况下,tasks数组里每个task的state都是RUNNING,你会看到类似这样的状态结果:
{ "name": "mm2-migrate-src-to-dst", "type": "source", "tasks": [ {"id": 0, "state": "RUNNING"}, {"id": 1, "state": "RUNNING"} ] }只要有task一直报FAILED,就得进日志查原因,最常见的几个原因后面单独开一节讲。
3.3 第三步:验证数据同步和位点同步
同步启动后,不要急着切流量。先在目标集群上确认几个关键现象:
先用topic列表确认目标集群上自动创建了哪些topic:
bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 --list正常情况下你会看到:
- 所有源集群的业务topic都出现了,并且topic名不变(因为DefaultReplicationPolicy不会加前缀)。
src.heartbeats、src.checkpoints.internal、mm2-offset-syncs.src.dst这类内部topic也出现了。
然后随机挑一个业务topic,分别查看源集群和目标集群的消息总量和最近offset对比:
bin/kafka-run-class.sh kafka.tools.GetOffsetShell \ --broker-list 192.168.1.10:9092 --topic order_events bin/kafka-run-class.sh kafka.tools.GetOffsetShell \ --broker-list 192.168.2.10:9092 --topic order_events注意,两个集群的offset数值会不一样是正常的,MM2在目标集群写入消息时,offset是目标集群自己分配的新offset,和源集群的offset没有可比性。我们要对比的是消息内容(消息数/最后一条消息的时间戳),不是offset数值。要验证数据完整性,更可靠的姿势是记录源集群每个分区的最新offset,再对比目标集群对应的每个分区消息数是否一致,或者直接在目标集群用消费者从开头消费一遍,抽样检查几条关键消息的内容。
位点同步的验证稍微复杂一点。先看源集群某个消费组的状态:
bin/kafka-consumer-groups.sh --bootstrap-server 192.168.1.10:9092 --describe \ --group order-service-group记录下当前的Lag情况,然后等待几个checkpoint周期(这个例子是5秒一个周期),在目标集群查同一个消费组的位点。你可以直接用--describe看目标集群的这个消费组是否已经有了offset记录,虽然一开始可能显示的是由mm2-...这个特殊group执行产生的记录,但确认存在且Lag在滚动更新就说明位点同步在跑。
如果这一步发现问题,最直接的手段是调小emit.checkpoints.interval.seconds,让快照更新得更频繁,然后在源集群手动消费几条消息,观察目标集群的位点是否跟着变化。
3.4 第四步:消费者切换的时机和操作细节
数据同步稳定运行一段时间后,业务方确认可以切换了。这里的标准流程是:
- 通知所有生产端停止写入或切到目标集群。先处理生产端,再处理消费端,这个顺序不能反。
- 等待存量数据全部同步完成。怎么判断?最简单的方法是在源集群查每个业务topic的LogEndOffset,再对比目标集群相同topic的LogEndOffset,等两者基本一致(目标集群不应再持续增长,或增长完全来自实时数据)。
- 让消费端切换broker地址到目标集群。消费组ID保持不变。
- 观察消费端的Lag和业务日志。正常情况下,消费者连接目标集群后,因为MM2已经通过checkpoints同步了位点,它会从大致对应的位置开始消费,你的业务代码看起来好像什么都没发生一样继续跑。
- 确认稳定运行后再做收尾清理。备份好的topic和消费组清单,确认旧集群不再有业务流量,再考虑下线MM2和旧集群。
关于第3步,有个很容易被问起的坑:如果消费端切过去后发现位点不对怎么办。这种情况通常是因为消费组在源集群上近期没有活跃消费,checkpoints没有及时更新,或者MM2的refresh.groups还没发现这个消费组。处理方式是用kafka-consumer-groups手动重置目标集群的消费位点,比如:
bin/kafka-consumer-groups.sh --bootstrap-server 192.168.2.10:9092 \ --group order-service-group \ --topic order_events \ --reset-offsets --to-datetime "2024-07-01T00:00:00.000" --execute重置到哪个时间点,取决于业务允许重复消费多少数据。如果一点都不能重复,那就得用--to-current或者精确offset,这个需要和业务方提前约定好。我的个人建议是:消费端能接受一定程度重复消费的话,切过去之前先重置offset到业务低峰期的时间点,让消费者在目标集群上从头消费low water或者最近几分钟的数据,这样即使MM2位点同步有秒级偏差,也不会丢消息,最多是复用几条老消息,影响很小。
3.5 迁移过程中的可视化监控方式
Kafka官方其实没有特别好用的MM2可视化界面,但我们可以借助一些开源工具来观察同步状态。热词里提到的kafka可视化工具,比如kafka-ui(现为UI for Apache Kafka)、Kafka Manager(雅虎开源,已停止活跃维护但能用),都可以同时配置多个集群地址,在同一个页面上切换查看源集群和目标集群的topic列表、消费组Lag。对迁移过程来说尤其方便的地方在于,你不需要在两套命令行之间来回切,直接对比两边的topic消息数和消费位点即可。
在kafka-ui里配置新集群的连接地址时,注意和MM2的目标集群保持一致,否则页面展示的数据和实际数据不一致会误导判断。另外,这类工具本身不建议部署在公司内网之外,它能看到所有topic的元数据信息,属于敏感资产,做好访问控制再暴露。我见过有人把kafka-ui放到公网服务器上方便查数据,这种操作等于把核心中间件元数据直接裸奔,尽量避免。
4. 常见问题与排查技巧实录:我踩过的那些坑
4.1 位点没同步:消费者切过去后从最新开始消费
现象:目标集群的topic数据是有了,但消费者切换后直接消费最新消息,老消息全被跳过。
原因:最常见的是MM2的refresh.groups还没发现该消费组,或者消费组在源集群上的位点信息还没被checkpoint捕获。第二个常见原因是目标集群上该消费组本身已经存在一个活跃消费实例,导致新consumer加入后重新分配分区并基于当前最新offset开始消费,覆盖了MM2同步过来的位点。
处理:
- 确认配置里
groups=.*,并且refresh.groups.interval.seconds设置合理(我一般用30秒,比5分钟更灵敏,代价是元数据请求频繁一点,但迁移期间能接受)。 - 切换前,先把目标集群上同名的消费组停掉,或者临时改个group.id让旧实例不占位,等MM2的checkpoint把源集群的位点覆盖过来之后,再恢复正式的group.id。
- 如果实在来不及,就直接用reset-offsets手动校准,不要犹豫。
4.2 同步延迟持续走高:复制进度跟不上生产速度
现象:源集群生产速率很高,目标集群的Lag一直涨,消费者切换过去之后永远消费不完积压的数据。
原因:tasks太少、checkpoint间隔太频繁、或者client端参数没调优。核心瓶颈通常在MM2的consumer拉取能力上。
处理:
- 先看Connect日志确认有没有rebalance频繁发生,如果task经常被重新分配,优先检查
tasks.max和topic分区数的关系。 - 调大MM2每条消息的拉取上限,在connector配置里加:
fetch.max.bytes决定一次拉取的总体大小,max.partition.fetch.bytes决定单分区拉取上限。"source.cluster.consumer.fetch.max.bytes": "104857600", "source.cluster.consumer.max.partition.fetch.bytes": "10485760" - 同时适当提高MM2生产端的吞吐,加:
注意,linger.ms增大意味着消息在生产者端多等一会儿才发出,会增加少量延迟,但能显著提升吞吐。迁移场景里有秒级延迟完全可以接受。"target.cluster.producer.batch.size": "1048576", "target.cluster.producer.linger.ms": "100" - 如果还不行,就需要对热topic做分队列处理:把最热的几个topic单独建一个connector,给它更高的priority和更多的tasks,让它和普通topic的复制任务隔离开,避免互相抢资源。
4.3 Connector反复失败并重启:日志里出现OffsetOutOfRange或TimeoutException
现象:Connect任务状态从RUNNING变FAILED,过一会又自动RUNNING,反反复复。
原因:常见的有两类——一是目标集群的某些内部topic(如checkpoints.topic)分区数不够,导致并发写冲突;二是MM2使用了过旧或过新的客户端协议访问目标集群,出现超时异常。
处理:
- 在启动Connect前,先手动创建内部topic,避免运行时自动创建带来的分区/副本参数不可控:
bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 \ --create --topic mm2-configs --partitions 1 --replication-factor 3 bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 \ --create --topic mm2-offsets --partitions 1 --replication-factor 3 bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 \ --create --topic mm2-status --partitions 1 --replication-factor 3 - 检查Connect日志(一般在
logs/connect.log),重点搜Caused by,很多问题都能从这里找到根因。如果日志级别不够详细,在connect-distributed.properties里把log4j.logger.org.apache.kafka.connect.mirror=DEBUG调上去,这个级别下每个topic的复制进度都会打印出来,排查延迟问题时尤其有用。 - 如果出现TimeoutException,检查两个集群的bootstrap.servers是否填错,以及Connect所在机器的防火墙是否放通了目标集群的9092端口。这个看似简单的问题,曾经让我排查了一整个下午。
4.4 有残留的内部topic和同步产物
现象:迁移完成、旧集群下线后,新集群里躺着一堆src.heartbeats、src.checkpoints.internal、mm2-offset-syncs.*。
原因:这是正常的,MM2在设计上就是通过topic来通信的,这些内部topic记录了迁移期间产生的元数据。但迁移结束后,它们就没有保留价值了。
处理:确认业务全部切到目标集群、源集群完全停用后,先删connectors:
curl -X DELETE http://192.168.1.20:8083/connectors/mm2-migrate-src-to-dst然后手动清理目标集群上的MM2内部topic:
bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 \ --delete --topic src.heartbeats bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 \ --delete --topic src.checkpoints.internal bin/kafka-topics.sh --bootstrap-server 192.168.2.10:9092 \ --delete --topic mm2-offset-syncs.src.dst这里特别强调一下,不要在生产集群上提前手动删除这些topic,否则正在运行中的MM2会一直被报错。内部topic的生命周期最好跟connector一致:connector活着,topic就留着;connector删了,再做清理。
4.5 常见问题排查速查表
| 现象 | 首要排查点 | 快速处理手段 |
|---|---|---|
| 目标集群topic没创建 | 检查connect日志、刷新间隔是否过了 | 手动执行一次refresh.topics或者重启connector |
| 目标集群有topic但消息数为0 | 检查正则是否有误、topic名是否被排除 | 用kafka-topics --describe确认配置同步状态 |
| 消费者切过去后重复消费大量消息 | checkpoints间隔过大或消费组未活跃 | 用reset-offsets校准到指定时间点 |
| 延迟持续上升 | tasks.max、fetch参数、热topic隔离 | 调大fetch参数,热topic拆独立connector |
| Sync失败且日志报ACL错误 | 源集群开了ACL但sync.topic.acls没开 | 开启sync.topic.acls.enabled=true并确保连接账号有读ACL权限 |
| Connect节点崩溃频繁 | 内存不足、堆外内存被元数据压爆 | 调大JVM堆内存,Connect默认堆只有1G,迁移场景至少给到4G以上 |
补充一个不那么起眼但坑过我的点:Connect的JVM堆内存。Kafka Connect默认的堆内存配置非常小,如果你迁移的topic特别多(几百上千个),元数据、client对象、内部缓冲都会吃堆内存,等到OOM就麻烦了。启动Connect前,改bin/kafka-run-class.sh里KAFKA_HEAP_OPTS到-Xmx6g -Xms6g,或者通过环境变量覆盖,这是最廉价又最有效的稳定性保障。
5. 收尾阶段的三个必做动作
5.1 在目标集群上验证数据完整性
这一步值得花时间做细。挑几个核心的黄金链路topic(比如下单、支付、库存这类关键业务topic),用生产者的同源数据做抽样比对外,还要确认分区数、副本数、retention配置、cleanup.policy和源集群一致。如果发现某个topic的配置在目标集群上不对,可以用kafka-configs --alter手动修正,或者把MM2里对应的topic排除掉,修正完再重新同步,注意不要影响其他topic。
5.2 打通生产端的双写校验
数据完整性验证通过后,要让业务生产端先切换到目标集群,并保持源集群的消费者继续运行一小段时间。这个阶段称为“灰度”。
为什么这么做?因为生产端切到目标集群后,数据开始只写新集群,源集群的数据不再增长。此时如果消费端还在源集群正常消费,那就说明消费逻辑没问题;等源集群的数据全部被消费完,消费端再切到目标集群,位点对不上。所以更稳妥的做法是,消费者也跟着生产端一起切,只是分批切,比如先切10%的消费实例到目标集群,让它们基于MM2同步的位点消费,观察一段时间没有消费异常,再逐步把剩余实例切过去。这个策略特别适合多副本、多实例部署的微服务架构,灰度窗口内新老集群同时在跑,出问题可以随时回切。
5.3 做好回退预案再下线旧集群
迁移做得再顺,回退的预案也要留着。正式切流量之前,把旧集群完整保留至少一个业务周期(我通常保留两周),中间不要急着清理任何topic和数据。一旦目标集群出现解决办法之外的严重问题(比如拉数据发现某个topic缺失、消费位点严重错乱),能把生产端和消费端地址直接指回旧集群,旧集群还在原来的位点继续服务,业务不至于中断。
回退这块还有一个容易忽略的点:旧集群的消费位点也是动态的。如果消费者切到目标集群后,源集群那边的消费实例并没有停,位点还在继续推进,那么切回去时消费位点可能跳跃很大。所以决定回退的话,要尽快在源集群上把所有消费组停住,然后也用reset-offsets校准一遍,别指望原来跑的位点还能直接接着用。
根据我自己的经验,迁移的成败很大程度不是看配置写得多对,而是看验证环节做得多细。上面这套流程我前前后后跑了不下十遍,每跑一次都会发现新的坑,比如alias命名不清晰、checkpoints同步时机不对、消费者提前占位导致位点被覆盖,这些问题在文档里基本不会有人提前提醒你。但只要你把准备清单列好、验证节奏踩稳、回退方案留着,MirrorMaker2这套方案在Kafka集群迁移里依然是最省心的那条路,毕竟它能把数据复制和位点同步这两件最繁琐的事情标准化掉,剩下的就看你的实操功夫了。