ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka跨集群迁移指南:MirrorMaker2配置、踩坑与消费切换实践

Kafka跨集群迁移指南:MirrorMaker2配置、踩坑与消费切换实践 做集群迁移这事最怕的不是数据量大而是迁移过程中业务还在跑存量数据和新写入的数据搅在一起割接时又发现消费端衔接不上最后变成一个没法收场的“半迁移”状态。我踩过这个坑之后遇到跨集群搬迁的需求第一反应就是整套MirrorMaker2的方案走一遍它有现成的复制通道、消费组位点同步机制、心跳检测和自动发现topic的能力虽然配置上有些细节容易踩雷但整体比自研迁移工具、或者用旧版MirrorMaker肉搏靠谱太多。这篇东西适合谁看手上有Kafka集群要换机房、升版本、拆集群的运维和开发还在纠结“双写 vs 迁移工具”怎么选的人以及第一次接触MirrorMaker2、想直接照着配置落地的同学。我会把从环境准备、配置设计、启动验证到消费者切换、问题排查的完整过程拆开写所有参数都给到能直接抄作业的程度并把我在实际操作中遇到的坑一并标注出来。1. 迁移方案选型为什么是MirrorMaker2而不是其他方案1.1 先盘一下市面上的主流迁移思路Kafka集群迁移常见的有三条路。第一是业务双写就是生产者同时往新旧集群发数据跑一段时间后切读。这个方案逻辑最简单但改动面很大所有生产端代码都要动而且双写期间消息顺序和幂等性都要额外保证对业务团队来说是实打实的侵入性改造。第二是自研复制工具自己写消费者拉旧集群数据、再生产到新集群好处是灵活坏处是offset管理、topic自动发现、分区映射这些全要自己处理稍不留神就出偏差时间成本也不低。第三就是用Kafka官方提供的MirrorMaker2MM2它本质上是Kafka Connect的一个连接器基于配置就能完成跨集群复制同时能同步消费组位点这是它相比自研方案最大的优势——不需要业务配合只在基础设施层操作就能把整个集群的topic和数据完完整整搬到新集群去。在真实的生产迁移场景里我基本不会考虑双写自研的方案除非公司有专门的中间件团队能长期维护一套迁移框架。大多数情况下业务方给到的迁移窗口只有几个小时这时候用MM2这种“开箱即用、可配置可监控”的方案才是性价比最高的选择。另外有些人会拿旧版MirrorMakerMM1来做对比这里必须强调一点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.heartbeatscheckpoints 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.idmm2-cluster # Connect实例的id每个节点唯一 client.idmm2-connect-1 # 存储connector配置、offset、状态的topic提前创建好或用auto.create config.storage.topicmm2-configs offset.storage.topicmm2-offsets status.storage.topicmm2-status # 这几个内部topic的副本数建议和集群保持一致 config.storage.replication.factor3 offset.storage.replication.factor3 status.storage.replication.factor3 # 指定目标集群或可同时访问两个集群的bootstrap地址 bootstrap.servers192.168.1.10:9092 # key/value converter建议用json调试更方便 key.converterorg.apache.kafka.connect.json.JsonConverter value.converterorg.apache.kafka.connect.json.JsonConverter key.converter.schemas.enablefalse value.converter.schemas.enablefalse # REST接口用来提交connector和管理任务 rest.port8083 rest.advertised.host.name192.168.1.20 rest.advertised.port8083这里有几个非常重要的点。第一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.max4这个值决定了数据复制任务拆成多少个task并行执行。MM2会按topic分区数自动做负载均衡task越多并行度越高但也不要无脑调高因为每个task都会占用独立的连接和线程资源建议先按topic分区总数除以100左右估算后续观察consumer组的Lag再调整。topics.*和groups.*用正则匹配所有topic和消费组。如果有特殊topic比如业务内部的高频日志topic不想要可以写成topics.*配合topics.excludeinternal_.*来过滤。refresh.topics.interval.seconds300每5分钟检查一次源集群有没有新增topic发现后自动在目标集群创建并开始复制。这个值太小会频繁请求源集群元数据太大会导致新增topic不能及时同步5分钟是合理折中。refresh.groups.interval.seconds300同步消费组信息的时间间隔同样影响位点同步的及时性。emit.checkpoints.interval.seconds5MM2每隔5秒把源集群消费组的offset快照写入checkpoints topic。这个值决定了迁移过程中目标集群能多快感知到源集群消费位点的变化。时间设得太大消费者切换时位点偏差就大。replication.factor3在目标集群创建topic时使用的副本数。这里很容易踩坑如果你把复制因子设得比目标集群的最小ISR还低那么同步过去后生产的可用性会出问题如果你设得比目标集群broker数还高topic创建直接失败。所以务必根据目标集群的节点数来定。sync.topic.configs.enabledtrue自动把源集群topic的配置同步到目标集群比如retention.ms、cleanup.policy这些这个功能很实用但要注意如果源集群有些topic配置比较特殊比如无限保留、超大分区同步过去后可能会影响目标集群的存储所以有特殊topic建议在topics.exclude里排除掉或者迁移后手动修正再放开同步。sync.topic.acls.enabledfalse如果源集群没开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列表确认目标集群上自动创建了哪些topicbin/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配置里加source.cluster.consumer.fetch.max.bytes: 104857600, source.cluster.consumer.max.partition.fetch.bytes: 10485760fetch.max.bytes决定一次拉取的总体大小max.partition.fetch.bytes决定单分区拉取上限。同时适当提高MM2生产端的吞吐加target.cluster.producer.batch.size: 1048576, target.cluster.producer.linger.ms: 100注意linger.ms增大意味着消息在生产者端多等一会儿才发出会增加少量延迟但能显著提升吞吐。迁移场景里有秒级延迟完全可以接受。如果还不行就需要对热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.mirrorDEBUG调上去这个级别下每个topic的复制进度都会打印出来排查延迟问题时尤其有用。如果出现TimeoutException检查两个集群的bootstrap.servers是否填错以及Connect所在机器的防火墙是否放通了目标集群的9092端口。这个看似简单的问题曾经让我排查了一整个下午。4.4 有残留的内部topic和同步产物现象迁移完成、旧集群下线后新集群里躺着一堆src.heartbeats、src.checkpoints.internal、mm2-offset-syncs.*。原因这是正常的MM2在设计上就是通过topic来通信的这些内部topic记录了迁移期间产生的元数据。但迁移结束后它们就没有保留价值了。处理确认业务全部切到目标集群、源集群完全停用后先删connectorscurl -X DELETE http://192.168.1.20:8083/connectors/mm2-migrate-src-to-dst然后手动清理目标集群上的MM2内部topicbin/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拆独立connectorSync失败且日志报ACL错误源集群开了ACL但sync.topic.acls没开开启sync.topic.acls.enabledtrue并确保连接账号有读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集群迁移里依然是最省心的那条路毕竟它能把数据复制和位点同步这两件最繁琐的事情标准化掉剩下的就看你的实操功夫了。
返回列表