
上个月刚完成一次Kafka集群数据迁移把一套用了快三年的集群从旧机房整体迁到了新环境。整个过程最花时间的地方不是复制消息本身而是怎么让业务无感切换。接手这个任务的第一个周末我把能查到的迁移方案都过了一遍最终选了Apache Kafka自带的MirrorMaker2下文简称MM2从验证、双跑、切换、清理四个阶段走完期间没有丢一条消息也没有出现长时间的业务中断。这篇把我个人落地这套集群迁移的完整步骤、配置文件和踩坑记录整理出来尤其是那些文档里不会写、实际特别容易绊倒人的细节。如果你正准备做Kafka集群跨机房搬迁、集群拆分合并或灾备建设这篇可以直接作为操作参考。1. 为什么选MirrorMaker2先盘点清楚其他方案的边界做迁移方案之前先要对场景做个分类。最常见的两类跨机房、跨环境整体搬迁集群拆分与合并比如把几个业务线的小集群合并成一个大集群。不同场景对“平滑程度”的要求不一样搬迁场景更多考虑切换时间窗合并场景更多考虑Topic归属和命名冲突。但无论哪类都有共同的硬指标不丢消息、不重复消费、切换快、能回滚。方案选型看的不是哪个工具炫酷而是它能不能满足这四个约束。我在这次迁移前列了一个需求清单旧集群不停服双跑至少一周新集群可随时回切Topic配置、ACL、消费组offset都要同步业务代码不改。有了这个清单自己写脚本转发的方案基本被否决了——数据复制容易消费位移和配置同步要自己实现工作量太大。1.1 其他方案的边界到底在哪很多人上来会问Kafka本身不是有副本机制吗为什么不能直接跨集群复制这里要分清一件事副本机制只在同一个集群内部的多broker之间同步数据跨集群复制跟它完全是两码事。kafka-reassign-partitions也不是干这个的它的职责是把同一个集群内的分区副本搬到别的broker上解决的是节点扩容不是集群迁移。常见的几个替代方案实际边界是这样的业务双写改业务代码新老集群同时写入。上线期间两端都要保证可用任何一端异常都会影响主链路切流时还容易乱序。除非新老系统需要长期并存否则不推荐用双写来做一次性迁移。旧版MirrorMakerMM1本质是消费源集群、写入目标集群能搬消息但Topic配置、ACL、消费组offset都不能同步。切换消费端时offset对不上必然出现重复消费或丢消息。自研消费转发程序要自己处理消费-生产的背压、幂等、配置同步、offset恢复还要考虑进程挂掉以后的断点续传非核心团队不建议碰。MM2Kafka 2.4.0之后内置底层基于Connect框架Source、Sink、Checkpoint、Heartbeat四类连接器分工明确能直接同步Topic配置和消费组offset正好覆盖迁移的全部需求。方案数据复制Topic配置/ACL同步消费offset同步改业务代码适合场景集群内副本扩容否否否否节点扩容业务双写是不涉及不涉及是新老系统长期并存MM1是否否否纯数据备份归档自研转发是需自研需自研一般不改有大量定制需求MM2是是是否跨集群迁移、灾备、合并选MM2还有一个现实理由它随Kafka版本一起发布不用额外维护一套独立组件升级路径也清晰。2. 迁移前盘点版本、Topic元数据与消费组依赖很多人拿到迁移任务就直接配MM2结果跑到一半发现目标集群的Topic配置不对、消息大小超限、ACL没同步。这些基本都能在迁移前盘点阶段提前暴露。2.1 版本匹配与broker参数核对MM2不是独立下载的软件它随Apache Kafka一起分发Kafka 2.4.0开始内置。我这次迁移是2.8.1到3.0.0两个版本中间有协议演进但MM2复制走的是Consumer和Producer的客户端路径兼容性风险相对小。不过还是建议源集群和目标集群的版本尽量一致至少跨度别超过一个大版本先做小流量验证再放开全量。有一个参数特别值得提前核对message.max.bytes和topic级别的max.message.bytes。源集群如果为某个业务调过1MB上限目标集群还是默认值复制过程中就会报RecordTooLargeException连接器任务反复失败重启。这类问题等全量复制时再发现排查成本很高。还有replica.fetch.max.bytes如果复制时单条消息较大、目标broker的这个参数偏小同样会出问题。稳妥起见把两端broker的这三个参数放到同一水平再从源端把topic级覆盖配置同步过去。2.2 Topic元数据盘点与预建建议复制Topic数据本身是自动的但不建议完全依赖MM2自动创建。源端Topic如果设置了cleanup.policycompact或很长的retention.ms目标端自动创建时可能用了默认的delete策略历史消息会被清理策略提前删掉。我这次的做法是迁移前把源端所有Topic的配置导出成清单和目标端逐一核对。盘点的字段至少要包含这些Topic名称分区数副本数cleanup.policydelete还是compactretention.msmax.message.bytes压缩类型compression.type分区数尤其要盯住。MM2复制时如果目标端Topic分区数小于源端可能出现消息分布不均衡的问题。建议在开启全量复制前直接把核心Topic在目标集群用kafka-topics.sh --create显式建好参数按源端配置对齐让MM2只负责数据同步不去依赖它的自动创建逻辑。2.3 认证、ACL与消费组依赖清单很多公司的Kafka开了SASL/SSL认证MM2连接器本身也需要凭据。要注意的是MM2开启sync.topic.acls.enabledtrue后会尝试把源集群的ACL同步到目标集群但两个集群用的如果是一套认证体系而principal不一致同步过去的ACL在鉴权系统里根本识别不了。这一点我后面会详细展开讲。消费组依赖清单也别漏。迁移前把所有消费组、订阅的Topic、当前Lag记录下来这些是后续切换和校验的基线。尤其要记下每个消费组的最新offset后面验证offset是否同步到位全靠它。bin/kafka-consumer-groups.sh --bootstrap-server 源集群:9092 --all-groups --describe输出里每个group的当前offset和Lag就是基线。建议保存到文件迁移结束后再跑一遍作对比。3. 吃透MM2复制机制alias、四类连接器与offset同步逻辑如果说上一节是准备这一节就是把MM2的关键机制讲透。我见过不少同事直接把官方配置粘过来就跑跑不通再回头查文档浪费很多时间。这里最值得先理解的就是alias机制。3.1 alias改名机制为什么目标集群会出现带前缀的TopicMM2会给每个集群配置一个alias别名默认的DefaultReplicationPolicy会把复制的Topic自动改成“目标集群别名.源Topic名”这种命名。举例源集群叫source目标集群叫target一个Topic叫order_event复制过去后你在目标集群看到的可能是target.order_event。这样设计的本意是防环双向复制时两条链路不会把同一条消息来回复制带上前缀就能区分消息来源。但这里有一个非常容易翻车的地方如果在目标端用业务消费者直接消费order_event会发现根本没有数据因为数据都在target.order_event里。很多第一次用MM2的团队都卡在这里排查半天。解决方式有三条路改replication.policy.class用IdentityReplicationPolicy保留Topic原名目标端业务消费的Topic名不用改。继续用DefaultReplicationPolicy但业务端消费带前缀的Topic相当于业务也要跟着改。在目标端再用一条复制链路把带前缀的Topic搬成原名但这样多一跳没必要。我这次选的是IdentityReplicationPolicy让Topic在目标端保持原名业务代码零改动。代价是如果后面做双向复制防环逻辑得自己处理好不能让消息在两侧无限循环。官方默认策略更安全但命名不透明Identity策略更透明但防环要靠自己。3.2 四个内置连接器的分工MM2不是单一组件它由四类连接器协同工作MirrorSourceConnector数据复制的核心消费源集群的Topic写入目标集群。复制过程中还会把消费位点写到专门的offset-syncs内部Topic。MirrorSinkConnector反向的Sink连接器配合Source使用。一般架构以Source为主Sink只在特定场景下才需要显式配置。MirrorCheckpointConnector定期把源集群消费组的offset“翻译”成目标集群消费组的offset写入checkpointTopic。这是消费者无感切换的关键。MirrorHeartbeatConnector周期性发送心跳消息记录端到端延迟用来观察两个集群的连通性和复制延迟。实战中启动MM2后你在Connect的REST接口里能看到这几个连接器自动被创建。如果你只看到Source连接器说明配置里没打开heartbeat和checkpoint那么后面消费者切换时会发现offset根本对不上。3.3 offset同步的机制与现实切换窗口offset同步是MM2相对MM1最大的优势但很多人对它有一个误解以为它是实时的。实际上源集群每个消费组都会定期提交offset到__consumer_offsetsMM2的Source连接器读到这些提交记录写入mm2-offsets内部TopicCheckpoint连接器再消费这个内部Topic把“这个消费组在目标端对应的Topic应该从哪个offset继续拉”写入目标集群的checkpointTopic目标集群的同名消费组从这个近似位置继续消费。注意“近似”这个词。复制链路本身是异步的checkpoint默认几分钟才生成一次所以当你把消费端切到目标集群时理论上会重复消费最近几分钟的数据。要完全避免需要先把消费端暂停等MM2完成最后一次checkpoint再切如果业务能做幂等直接把checkpoint周期调短并接受小窗口重复。这个点一定要在切换前跟业务方讲清楚异步复制方案都不存在零窗口提前沟通能避免很多问题。4. 生产级配置与启动验证从单Topic跑通到全量同步原理清楚了配置才有底气。下面是我这次迁移实际用到的配置文件模板关键参数都写了注释。4.1 mm2.properties逐段解读# 定义两个集群的别名 clusters source, target # 源集群和目标集群的broker地址 source.bootstrap.servers 10.0.0.1:9092,10.0.0.2:9092,10.0.0.3:9092 target.bootstrap.servers 10.0.0.11:9092,10.0.0.12:9092,10.0.0.13:9092 # 用IdentityReplicationPolicy保留Topic原名业务端不用改消费Topic replication.policy.classorg.apache.kafka.connect.mirror.IdentityReplicationPolicy # 全局默认任务数 tasks.max8 # Topic和消费组的刷新周期 refresh.topics.interval.seconds60 refresh.groups.interval.seconds60 # 开启从source到target的复制 source-target.enabledtrue # 复制所有Topic和所有消费组 # 第一次跑建议先用具体Topic名验证跑通后再放开成 .* source-target.topics.* source-target.groups.* # 开启心跳和checkpoint source-target.emit.heartbeats.enabledtrue source-target.emit.checkpoints.enabledtrue # 同步Topic配置和ACL source-target.sync.topic.configs.enabledtrue source-target.sync.topic.acls.enabledtrue # 内部Topic副本数务必根据目标集群broker数量调整 source-target.heartbeats.topic.replication.factor3 source-target.checkpoints.topic.replication.factor3 source-target.offset-syncs.topic.replication.factor3逐个说下关键点。clusters定义了别名后面所有配置都用source-target.前缀限定方向这个命名空间设计让双向复制变得清晰。topics.*看起来很省事但第一次跑最好不要直接上我是先用一个kafka-migration-check验证Topic跑通再重启补成.*避免配置错误把不想搬的内部Topic也搬过去。sync.topic.acls.enabledtrue这个参数生产环境我建议先确认两个集群的principal体系一致再开否则同步过去的ACL识别不了等于白开。heartbeats.topic.replication.factor和checkpoints.topic.replication.factor要根据目标集群broker数量定如果目标集群是3个broker设成3没问题如果是2个broker设成3会直接报错。4.2 启动方式选择与健康检查MM2支持单机standalone和分布式两种模式。生产一定用分布式因为分布式模式会把连接器状态存到Connect内部Topic里支持横向加节点进程挂了能自动恢复任务分配。启动命令比较简单较新版本一般直接这样起bin/connect-mirror-maker.sh config/connect-mirror-maker.properties如果集群开了认证还需要在启动前把两个集群的JAAS配置都写清楚不要把生产端和消费端凭据混在一起。这一步最容易漏等连接器全部FAILED再回头查配置很耽误事。启动后第一件事是看连接器状态curl localhost:8083/connectors curl localhost:8083/connectors/连接器名/status期望看到state是RUNNINGtasks里每个task也都是RUNNING。如果出现FAILED一般是认证、Topic副本因子、消息大小这几类问题。这时候去看Connect的日志通常第一行报错信息就能定位。4.3 用一个小Topic验证全链路正式全量之前强烈建议先用一个小Topic验证全链路。具体做法建一个分区数为3、副本数为2的验证Topic往源集群写入几百条标记消息同时启动一个简单的消费组消费它。等MM2把数据复制过去在目标集群确认Topic分区数一致、消费组offset被同步这个小实验就算通过。# 源集群写入消息 bin/kafka-console-producer.sh --bootstrap-server 源集群:9092 --topic kafka-migration-check # 目标集群消费消息 bin/kafka-console-consumer.sh --bootstrap-server 目标集群:9092 --topic kafka-migration-check --from-beginning对比消费组offset时用同一条命令在两个集群分别跑bin/kafka-consumer-groups.sh --bootstrap-server 目标集群:9092 --group 某group --describe如果目标集群能看到同名group且Lag值在缩小说明offset同步链路是通的。这个小实验大概花半小时能挡住90%以上的配置错误。5. 双跑验证、流量切换与回滚底线我的实操顺序配置跑通不代表能直接切从“复制正常”到“可以切换”之间还隔着数据校验和切换顺序设计。这一节分享我的实操顺序尤其切换顺序我踩过对比坑之后才确定下来。5.1 数据一致性校验总量对比与抽样确认数据迁移最怕复制了一半没发现。我的校验方法分两层。第一层是总量对比。分别统计源和目标所有Topic的分区消息总数diff为0说明大体复制成功。但双跑期间两个集群不断有新数据进来总量diff永远不为0所以实操中我把对比限定在某个时间点之前的数据先让全量追平再停止写入一小段时间做一次快照对比。bin/kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list 源集群:9092 --topic order_event --time -1 bin/kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list 目标集群:9092 --topic order_event --time -1第二层是抽样确认。对核心Topic随机挑几个分区读取该分区最后N条消息在生产端造几条带唯一标记的消息在目标集群用消费者把这些标记消息捞出来确认端到端链路延迟在可接受范围。MM2的心跳Topic是最直接的延迟观测手段源端发一条心跳目标端何时收到两者时间差就是复制延迟。这个指标是判断能否切换的重要参考。5.2 切换顺序为什么我坚持消费者先切关于切换顺序网上有两种声音先切生产者还是先切消费者。我试过之后强烈建议消费者先切、生产者后切而且中间隔至少一两个checkpoint周期。理由很简单。消费者切到目标集群后如果offset同步有窗口最多重复消费一小段幂等设计完善的业务应该能接受。如果先切生产者旧集群消费端将无法消费新数据一出问题整个业务就断了。双跑期间旧集群继续接收全部流量新集群一直通过MM2追平消费者切过去之后再等目标集群Lag清零此时切生产者。即使生产者切过去之后出现问题目标集群仍然有完整数据回滚也只需要把生产者和消费者都指回旧集群。这个顺序最大的好处是任何时刻出错业务都还能回退到旧集群不会出现两端都不可用的情况。5.3 切换操作清单与回滚预案我整理的切换清单发布操作时直接照着勾选源集群开启全量Topic复制确认所有核心Topic已在目标集群存在。暂停非核心消费组减少迁移期间的干扰。观察MM2心跳延迟连续3个周期延迟低于阈值。逐个服务把消费者bootstrap.servers切到目标集群发布后观察消费组Lag。等待至少两个checkpoint周期确认所有消费组在新集群恢复消费。逐个服务把生产者bootstrap.servers切到目标集群。记录切换时间点作为后续数据一致性检查的基准线。保留旧集群运行观察一周后再决定下线。回滚预案的核心是只要旧集群还在运行任何时刻都可以把生产者和消费者的bootstrap.servers指回旧集群。要点是在回滚前确认旧集群没有被MM2反向同步单向复制配置下不会出现这个问题。6. 迁移实录六个容易翻车的细节与根因这一节把这次迁移中真实踩过的坑写出来每个都附上根因和解决方式。这些坑在官方文档里很难一次性看全但实际工作中几乎都会遇到。6.1 认证模式下ACL未同步目标集群消费端权限全挂我这次用的是SASL_PLAINTEXT加ACL授权MM2配置里写了sync.topic.acls.enabledtrue以为目标集群的Topic权限会自动建好。实际上目标集群的KafkaPrincipal跟源集群不完全一致MM2同步过去的ACL在鉴权系统里没有被识别。排查过程目标集群消费者一直报TOPIC_AUTHORIZATION_FAILED我先查了目标端Topic ACL发现授权记录是有的但记录的principal和实际连接用的principal对不上。最后手动重建了目标集群的业务账号和ACL再关掉自动ACL同步避免覆盖问题才解决。教训是自动ACL同步看起来很美好生产环境还是建议目标集群提前手动配好ACLMM2只负责同步数据。6.2 消息大小限制不一致大消息复制直接失败这个坑跟很多人搜的kafka 接收1m直接相关。源集群原来为某个业务调过message.max.bytes1048576目标集群还是默认值复制过程中连接器任务报RecordTooLargeException任务反复失败重启。我一度以为是MM2版本问题查了很久才发现是两边参数不一致。最后统一了两边的message.max.bytes、replica.fetch.max.bytes同时把源端Topic的覆盖配置同步到目标端问题才消失。这个坑特别隐蔽因为源端单个消息不超1MB平时根本不会触发一旦业务上传大对象就炸。6.3 checkpoint周期内的重复消费窗口前文提过用MM2做迁移时offset同步是按checkpoint周期进行的默认5分钟一次。消费者切到目标集群后可能从源端5分钟前的位点开始消费重复消费窗口必然存在。如果你的业务完全无法容忍重复可以在切消费者之前先把源端消费组暂停等MM2完成最后一次checkpoint再切。如果业务能做幂等直接把source-target.emit.checkpoints.interval.seconds调短到30秒同时接受小窗口重复。这个参数的调整要在切流前完成切过去再调就晚了。6.4 双向复制不做防环消息在两个集群之间打转如果要做双活或多活两个集群都运行MM2且用了IdentityReplicationPolicy又没有在Topic正则里排除内部Topic和环回消息一条消息可能被A复制到B又被B复制回A如此反复。MM2官方的DefaultReplicationPolicy能自动区分来源但Identity策略不去区分只能靠自己在正则里排除或自定义策略。我的建议是双向场景下要么用默认策略让消息带来源前缀要么自定义ReplicationPolicy标记来源集群。内部Topic以mm2-、heartbeats、checkpoints开头的一定要在复制正则里排除避免无限膨胀。6.5 连接器任务分配不均单broker流量打满MM2的tasks.max和Topic分区数需要匹配。如果Topic有100个分区但tasks.max设成8连接器会尽量均衡但少数大Topic的复制仍会集中在某个task上迁移高峰期源集群某个broker可能出现超高读流量。我后来对大Topic单独开了一组MM2实例用独立的mm2.properties启动多个进程或者在一个配置文件里用不同的连接器配置加上topics.include和topics.exclude分组。核心大Topic独立连接器实例后其他Topic的复制延迟立刻降了下来。6.6 目标集群磁盘容量估算太乐观差点撑爆迁移新集群如果不能复用老集群的保留策略默认retention7天可能导致目标集群磁盘压力显著高于源集群。如果源端本来就是7天保留新增副本加上心跳、checkpoint内部Topic也会占空间。迁移前可以用一个简单公式估算每日写入量×保留天数×副本因子×1.3余量。别以为同规格迁移就没问题实际内部Topic数量比源集群多了不少。我这边就是因为低估了这部分迁移第三天才发现某台broker的磁盘使用率涨得比预期快紧急调了retention才压住。7. 迁移后的观察清单用一周时间确认新集群稳定切换完成不代表事情结束。前一周是问题集中暴露的时期我一般盯四类指标消费组Lag曲线、MM2心跳延迟、集群磁盘和ISR状态、业务侧错误率。尤其注意Lag是否持续增长或抖动磁盘是否因为复制积压而上涨过快。内部Topicmm2-开头如果长期只增不减说明需要给它们设置单独的retention.ms不然会一直占用磁盘。7.1 一周内重点观察的指标消费组Lag是第一个要盯的。切换后如果某个消费组在新集群上的Lag持续不降优先查该组用的Topic在目标集群的分区数是不是和源集群一致以及MM2的checkpoint周期是否覆盖了它的offset更新频率。MM2心跳延迟如果突然变大大概率是某个复制任务卡住了这时候去Connect的REST接口看task状态最直接。磁盘和ISR状态属于基础监控。切换后新集群的流量结构发生变化磁盘增长曲线如果连续两天斜率异常要尽快排查是否有Topic的retention配置没同步过来。7.2 检查命令与预期结果汇总下面是我迁移后每天跑一遍的检查项可以直接参考检查项命令/工具预期结果Topic数量kafka-topics.sh --list目标集群与源集群一致分区副本kafka-topics.sh --describe分区数一致副本数达标消费Lagkafka-consumer-groups.sh --describe所有group稳定无持续增长复制延迟目标集群消费heartbeat Topic延迟稳定在秒级连接器状态curl /connectors/name/status全部RUNNING我自己做这类迁移的最大体会是工具能解决数据复制的问题但真正决定迁移成不成的是迁移前对业务消费链路的盘点以及迁移后一周的耐心观察。MM2的每一条内置连接器都对应一类问题提前把原理吃透配置只是抄作业的事。如果你的场景跟我不一样比如要做跨地域双活那防环和冲突处理还要单独设计。希望这篇踩坑记录能让你少走几步弯路。