
刚完成一次Kafka集群迁移整个过程踩了不少坑也积累了一些经验。正好有朋友问起MirrorMaker2的具体用法我就把这次迁移的方案设计、配置细节、实操过程以及遇到的各种问题整理成文希望能给准备做Kafka跨集群复制的朋友一些参考。如果你有Kafka使用基础正在为集群迁移、数据同步、异地容灾发愁这篇文章应该能帮到你。1. 迁移方案的选型与整体设计1.1 先搞清楚为什么做迁移迁移前要确认哪些事先说我们这次迁移的背景。我负责的业务线有两个自建Kafka集群一个承载线上核心业务消息另一个主要是日志采集和异步任务。随着数据量涨上来旧集群的磁盘、网络和控制器节点都开始吃紧同时Kafka版本还停留在旧版本很多新版特性没法用安全和性能都有隐患。所以团队打算把两个集群平滑迁移到一套新的、规格更高的Kafka集群上同时把版本一次升到位。做技术方案之前我建议先想清楚这几点业务对停写时间的容忍度。有些业务能接受凌晨短暂停写有些核心链路要求迁移过程完全无感。这直接决定你能用多“暴力”的迁移手段。迁移的数据规模。总消息量多少、Topic数量和分区数多少、单条消息大小分布如何这些决定了同步带宽和耗时预估。消费组数量和消费位点状态。迁移不只是搬数据还要搬运“消费进度”消费者换集群之后能不能接着上次的进度消费这是最容易被忽略的。上下游依赖。有没有服务在往Kafka里面写数据、有没有实时计算任务在消费、数据有没有外部归档链路迁移时这些都是变量。把这些信息收集齐了再谈方案选型不然很容易做到一半发现某个系统写死了旧集群地址又得回去改。1.2 主流迁移方案对比为什么最终选MirrorMaker2Kafka跨集群复制和数据迁移的常用方案其实就几类我简单列一下各自的优缺点方案优点缺点适用场景业务代码双写实现简单、直接控制两条链路需要改动业务代码、存在写入不一致风险数据量小、能接受代码改造停机导数据 重置消费位点思路直观、工具生态成熟停机时间长、需要人工核对offset允许长时间停写原生跨集群复制方案生态配套比较完善配置复杂、资源占用高有一定Kafka运维能力的团队MirrorMaker2无需改业务代码、配置相对简单、自带消费位点同步对连接器稳定性要求高绝大多数跨集群复制场景我们最后选了MirrorMaker2核心原因是它不需要改业务代码也不要求消费者端做特殊改造。它本质上是构建在Kafka Connect框架上的一组连接器能把源集群的数据复制到目标集群同时把消费组的位点信息也一并同步过去消费者切到新集群之后用同一个消费组名接着消费基本能做到无感切换。这一点是手动重置位点方案很难做到的。补充一句这里说的“无需改业务代码”是指数据同步链路不需要业务配合但你的服务配置里Kafka地址那一项肯定还是要改的这是绕不开的。1.3 整体架构设计集群拓扑、Topic映射与同步模式迁移架构上我们采用了单向同步也就是把旧集群的数据复制到新集群。源集群和目标集群的拓扑关系我画了个简单示意源集群A承载原有全部Topic版本为旧版。目标集群B新的Kafka集群版本升级到新的稳定版Topic在同步开始后自动创建。MirrorMaker2节点独立部署不跟业务混部避免互相影响。这里有一个非常关键的设计点Topic映射关系。MirrorMaker2默认的复制策略会把源集群的Topic复制到目标集群后改名规则是“源集群别名.Topic名”比如源集群别名是A那么Topic为“order_event”的数据到目标集群就变成“A.order_event”。这样设计是为了支持双机房双向复制时避免循环复制。但我们迁移场景希望Topic名称保持一致不然业务切换后所有Topic名都得改配置文件、监控告警、下游任务全都要跟着动。解决办法是使用IdentityReplicationPolicy来保持Topic名不变。这个选择我们在后面配置部分会具体说明。同步模式上我选择的是单向复制。迁移期间让两个集群暂时并行运行消费者逐渐切到新集群等新集群数据追平、业务稳定后再停掉同步任务。这是比较稳妥的“并行运行-逐步切换”思路回滚也方便。2. MirrorMaker2核心原理与配置项拆解2.1 三个内置连接器各自在干什么刚接触MirrorMaker2时我一度以为它就是另一个“把消息从A搬到B的工具”直到在线排查问题时才理清它是一组连接器协同工作的结果。这里非常有必要把它的三个内置连接器讲清楚。MirrorMaker2在启动时会同时创建三类连接器MirrorSourceConnector负责把源集群的Topic数据复制到目标集群同时会把源集群的分区数、副本数、配置同步到目标集群。数据复制的主体是它。MirrorCheckpointConnector负责同步消费组位点。它会周期性地把源集群各消费组在每个分区上的offset记录下来写到目标集群的内部Topic消费者切到目标集群后按消费组名就能从这里拿回位点继续消费。MirrorHeartbeatConnector负责发送心跳消息。它定期往目标集群写一条特殊格式的Topic默认叫heartbeats用来验证源集群和目标集群的连通性以及复制链路是否健康。初次上手的人容易忽略后面两个。其实它们才是“无感切换”的关键。没有MirrorCheckpointConnector消费者切到新集群后只能从最新位置开始消费或者你自己手动去翻译offset非常痛苦。没有Heartbeat你连复制链路是否存活都只能靠翻日志。2.2 配置文件逐项拆解理解了三个连接器的职责配置就顺理成章了。MirrorMaker2的配置文件实际上就是一组Kafka Connect的配置但多了几个集群别名和同步策略相关的参数。下面是我们迁移时使用的一份完整配置示例# 集群别名逗号分隔 clusters A, B # 指定每个集群的broker地址 A.bootstrap.servers 10.0.0.1:9092 B.bootstrap.servers 10.0.1.1:9092 # 连接器使用到的存储Topic配置 A.offset.storage.topic mm2-offsets A.status.storage.topic mm2-status A.config.storage.topic mm2-config B.offset.storage.topic mm2-offsets-b B.status.storage.topic mm2-status-b B.config.storage.topic mm2-config-b # 源集群和目标集群的别名定义 A.alias A B.alias B # 指定要同步的Topic支持正则 topics .* # 不同步的内部Topic topics.exclude mm2-offsets.*|mm2-status.*|mm2-config.*|heartbeats|A.heartbeats # 使用保持Topic名称不变的策略迁移场景的关键 replication.policy.class org.apache.kafka.connect.mirror.IdentityReplicationPolicy # 同步Topic的副本因子 replication.factor 3 # 检查配置变更、Topic变更的间隔 refresh.topics.interval.seconds 60 refresh.topics.enabled true # 同步Topic ACL配置的开关 sync.topic.acls.enabled false # 心跳相关配置 heartbeat.interval.ms 5000 emit.heartbeats.interval.seconds 5 # 位点同步相关配置 checkpoint.checkpoint.interval.ms 10000 emit.checkpoints.enabled true # 连接器Task数量看情况调整 tasks.max 4 key.converter org.apache.kafka.connect.converters.ByteArrayConverter value.converter org.apache.kafka.connect.converters.ByteArrayConverter逐个说几个关键参数这些是我实际踩过坑之后才真正理解的clusters集群别名非常重要同时要求配置文件的后续前缀与之一致。比如A.bootstrap.servers中A是源于clusters声明的别名不能凭空定义。别名的选择会影响Topic的重命名策略如果用的是默认策略我们这里用A和B配合IdentityReplicationPolicy后不受影响。topics .*把源集群所有Topic都同步过去。如果你只想迁移某些业务Topic这里可以写具体的Topic列表或正则比如topic-a|topic-b。replication.policy.class这里值得再强调一次。默认是DefaultReplicationPolicy会把Topic改名成“源集群别名.Topic名”。我们用IdentityReplicationPolicy保持原名迁移场景基本都用这个避免下游配置和监控大面积修改。topics.exclude需要把MM2自身产生的内部Topic和心跳Topic排除掉否则会造成内部数据无限循环复制。表达式里我加了heartbeats|A.heartbeats就是排掉心跳Topic。key.converter和value.converter使用ByteArrayConverter默认的JsonConverter会尝试反序列化消息内容遇到非JSON数据会直接同步失败。Kafka消息对Kafka本身来说就是字节数组所以直接用字节转换器最安全原样搬不做任何解析。注意我这里给的配置只展示了核心项。实际生产环境还需要根据集群Soft的实际情况设置status.storage.replication.factor、offset.storage.replication.factor、config.storage.replication.factor等默认值可能不够安全数据量大的时候内部Topic的副本数也要保障。2.3 内部Topic和数据流MirrorMaker2把状态放哪里了MirrorMaker2运行时会依赖一组内部Topic理解这些Topic能帮你快速定位复制链路问题。初次排查时我看到一堆带mm2前缀的Topic很懵顺手整理了一张表内部Topic存什么数据主要用途mm2-offsets源集群各消费组的消费位点Checkpoint连接器读写位点mm2-status连接器运行状态Connect框架内部状态mm2-config连接器配置信息Connect框架内部配置heartbeats心跳消息检查集群连通性A.checkpoints.internal消费组位点快照目标集群消费组起始位点来源mm2-config.B.A同步策略产生的配置内部管理这些Topic一般都会自动创建。在权限严格的环境里如果你的Kafka开启了ACL需要提前给MirrorMaker2账号授予对这些内部Topic的读写权限否则启动时会一直报错看起来像“网络不通”其实是权限不足。数据流方向是MirrorSourceConnector自动从源集群拉取消息写入目标集群MirrorCheckpointConnector从源集群读取消费组位点然后写入目标集群的checkpoints内部Topic消费者切换后在目标集群以相同消费组名订阅Topic就能通过内部Topic找回位点。理解了这个链路后续排查“为什么切过去之后消费位置不对”这类问题就会清晰很多。3. 迁移实操实录与切换细节3.1 环境梳理与启动准备先说下沉环境。我们用了独立的机器部署MirrorMaker2没有和业务服务或Kafka Broker混部。原因很简单MirrorMaker2本质是Kafka Connect它在同步数据时需要占用相当多的网络带宽和内存如果跟其他服务混部要么影响业务要么影响同步速度两头吃力。部署包选择上关键点在于Kafka版本的兼容性。当时我们源集群是旧版本目标集群打算升到新版。官方文档说MirrorMaker2从某个版本开始才具备较完整的功能保险做法是把MirrorMaker2的版本升级到目标集群版本然后同时连接源集群和目标集群。新旧版本的Kafka协议兼容是没问题的这样既能充分利用新版本连接器的能力和修复又能连接旧版本集群。这一点不少朋友配置完启动报各种奇怪的错大概率是用了一个和源集群同版本的旧包而那个版本对某些新语法支持不全。启动之前要准备好三块东西配置文件也就是上一节那份单独放到config/目录下。从Kafka安装目录里找到connect-mirror-maker.sh脚本如果安装包完整这个脚本在bin/目录下。为MirrorMaker2创建一个独立的系统服务脚本用systemd或supervisor管理都行这样宕机之后能自动拉起避免人工介入。启动命令很简单bin/connect-mirror-maker.sh config/mm2.properties但启动之前强烈建议先在命令行前台跑一次看日志是否正常再交给守护进程管理。因为配置错误直接后台化之后排查起来很别扭。3.2 从全量同步到增量追赶怎么看同步进度启动MirrorMaker2之后它会立即开始把源集群所有匹配的Topic数据复制到目标集群。这个阶段相当于“全量同步”目标集群没有数据所以会以较快速度追平历史消息。数据量大的时候这个过程可能持续几小时甚至一两天取决于消息总量和带宽。同步过程中怎么判断进度我一般看以下几个指标目标集群的Topic是否都建起来了。执行kafka-topics.sh --bootstrap-server 目标集群地址 --list看Topic数量和源集群是否一致。同步过程中Topic会自动创建如果缺了某个Topic多半是topics正则没匹配到或者权限有问题。每个Topic的各分区offset是否接近源集群。用kafka-run-class.sh kafka.tools.GetOffsetShell分别查源和目标集群的同一Topic各分区的最新区位点对比差值。差值收敛到秒级甚至毫秒级说明增量追赶基本完成。查看MirrorMaker2监控指标。Kafka Connect暴露了JMX指标里面有复制延迟相关的指标项。通过JMX工具或者普罗米修斯采集都能看到主指标就是各Topic分区的复制滞后量。这个阶段最常犯的错误是只看Topic数量对上了就以为同步完了。其实Topic数量对上只说明结构同步完成数据不一定追平一定要对比分区offset。我们有次就是Topic都建好了结果一查某个核心Topic还有一千多万条消息的滞后直接切流量的话数据会丢一大截。3.3 与业务约定切换窗口平滑切换与回滚设计等数据追平之后才进入真正的“切换”环节。这个环节要和业务方、下游消费方提前约定时间窗口一般来说是凌晨或业务低峰期。我们的切换流程分这几步停掉源集群的消费方或者把消费方切到目标集群读取。确认MirrorCheckpointConnector已经把最新位点写过来了在目标集群上查看checkpoints相关内部Topic是否有新数据确认位点更新时间在最近一分钟内。把生产者的写入地址切到目标集群让新数据直接写入新集群。观察目标集群的生产消费是否正常、消费Lag是否逐步下降、日志有无报错。持续观察一段时间后确认稳定再停掉MirrorMaker2同步任务。这里要特别强调一个容易踩的坑先切消费者还是先切生产者是有讲究的。我建议先切消费者再切生产者。因为如果先切生产者新数据直接写到目标集群但部分消费者还没切过去那这部分消费者在源集群上就再也读不到刚产生的新消息了。反过来先切消费者让消费者从目标集群读这时候MirrorMaker2还在同步新旧数据都在目标集群消费者不会漏消息。回滚方案同样重要。我们在切换时保留了源集群的写入链路一周一旦目标集群出现重大问题可以把生产者和消费者都切回源集群。因为MirrorMaker2是单向同步的源集群的数据不会受目标集群影响回滚只需要改配置重启服务就行不需要重新导数据。注意回滚窗口不是无限的。如果切换后目标集群跑了两三天源集群那边可能已经有数据被清理策略删掉了此时再回滚就会出现数据缺失。所以确定新集群稳定之后要尽快释放源集群资源不要留太久。3.4 切换后的验证清单切换完成并不代表迁移结束我还整理了一份验证清单全部通过才算真正完事核心Topic在目标集群上的消息总量与源集群迁移前基本一致允许增量差异。各消费组在目标集群上的committed offset与源集群切换前的offset对应。新写入的消息能在目标集群被生产、消费正常。监控面板没有出现持续的复制异常、权限报错、连接超时。下游依赖Kafka的计算任务没有出现重复消费或大量积压。这套清单看起来繁琐但能帮你在切换后的混乱期快速定位问题而不是逐个人问“你那边正常吗”。4. 常见问题排查与避坑记录4.1 同步性能瓶颈背压、线程与内存配置跨集群复制看起来简单线性扩展其实有不少讲究。第一次实操时我把tasks.max设成4觉得够用了结果同步单个大Topic时延迟一直在涨。查了半天发现瓶颈不在CPU而在默认的发送端缓冲和内存上。MirrorMaker2同步性能受几个因素影响tasks.max数量它决定了连接器任务并行度适当调大可以提高吞吐。但并不是越大越好每个Task都会占用内存和连接资源集群Broker也会受影响。record累加器的内存Kafka生产者端的buffer.memory默认配置在低配机器上可能不够大数据量时会导致频繁阻塞。可以适当调大。端到端延迟通过两个集群的offset差值能反映出来如果差值持续扩大优先看网络带宽和生产者发送性能而不是盲目加Task。内存这块我建议给MirrorMaker2分配足够的堆内存同时在启动脚本里显式设置Kafka客户端的内存相关参数。JMX指标里有一个非常关键的指标是“复制延迟”建议配到监控告警里一旦延迟超过阈值就报警避免问题发酵成大事故。有一个容易被忽视的点Kafka Connect框架本身有offset.flush.interval.ms这类参数如果设置过小会导致频繁刷写内部状态反而降低吞吐。我们当时没动这些参数用默认值是可以的但需要知道它们的存在。4.2 配置不自洽引发的疑难杂症MirrorMaker2的报错信息有时候相当迷惑我整理了这次实操中遇到的高频问题的排查表现象常见原因快速排查方式目标集群一直不创建Topictopics配置写错、正则不匹配检查配置里的topics项先只同步一个测试Topic验证启动后大量Reassignment或Timeout日志内部Topic分区数过多、Broker负载高检查mm2-开头内部Topic的分区数必要时清空重建同步的Topic改名成了A.xxx未使用IdentityReplicationPolicy检查replication.policy.class配置消费者切过去后从最新消息开始读MirrorCheckpointConnector没生效或emit.checkpoints.enabled设成了false查看checkpoints内部Topic是否有数据同步速度慢、经常报网络超时Broker连接数不足或者带宽不够用ss看连接数检查机器带宽占用日志报ACL权限错误sync.topic.acls.enabled开启但目标集群权限配置不完整迁移场景通常直接改为false其中“目标集群Topic改名”是最容易发现的坑但也是最容易忽略后果的坑。如果已经用了默认策略跑起来了改配置之后需要重启连接器已经创建的错误Topic需要手动删掉重新同步。4.3 流控与稳定性别让同步任务拖垮业务最后聊一个全局性的经验。MirrorMaker2同步数据时源集群和目标集群都可能承载大量额外的读写压力。如果源集群本身是核心业务集群同步流量一定要控制住。怎么控制几个方式在MirrorMaker2的Topic同步列表里把非必要的Topic先排除掉。批量迁移时先同步核心业务Topic日志类低优先级Topic放到后面再开降低初期带宽占用。给MirrorMaker2的客户端设置适当的流控参数。Kafka客户端本身提供了producer端的限流参数按需配置可以平滑发送速率。错峰同步。如果业务有明显的波峰波谷可以把高流量Topic的同步节奏错开到低峰期比如通过分批修改topics正则来实现。我个人体会最深的一条是迁移方案设计的重心从来不是“怎么把数据搬过去”而是“怎么保证切换过程中业务无感、可控、可回滚”。MirrorMaker2只是把数据搬运这个脏活累活做好了剩下的切换编排、验证和回滚还是要靠人来做。这次迁移项目我们前后花了一周多时间真正执行切换只用了不到半小时。能在切换环节这么顺利主要还是前面把同步位点、Topic命名、权限策略都提前验证透了。希望这篇分享能让你在Kafka集群迁移时少走一些弯路尤其是那些只在生产环境才会遇到的坑提前避开能省下不少半夜排查故障的时间。