☰
Kafka迁移实战:场景判断、方案选型与MirrorMaker2落地指南
2026/10/8 20:12:37 网站建设 项目流程

最近好几个朋友都不约而同地问我同一个问题:Kafka迁移到底应该怎么做?有人要把自建集群搬到云上,有人要把旧版本集群升级到新架构,还有人只是想换掉一批老节点,结果发现网上那些操作手册和自己遇到的场景完全对不上。说实话,Kafka迁移这个题目看着简单,真做起来坑非常多。它不是一个“把数据从A拷贝到B”的问题,而是一个需要同时考虑数据一致性、消费者位点、Topic配置兼容、参数对齐、切换窗口、回滚预案的系统工程。这篇文章我就把做过几次迁移的经验完整梳理一遍,从场景判断、方案选型到实际操作、问题排查都覆盖到,适合正在搞集群搬迁、Kafka升级、节点替换的运维和开发同学。

1. 迁移前先想清楚:你是哪种迁移,要解决什么问题

1.1 三种典型迁移场景:集群搬迁、节点扩容/缩容、跨版本升级

Kafka迁移的第一步不是敲命令,而是先给这次迁移定性。我见过太多人上来就搜“MirrorMaker怎么配置”,结果他的场景根本不需要跨集群复制,白白搭进去一个组件还有一堆运维负担。按我的经验,日常遇到的迁移基本可以归成三类。

第一类是整体集群搬迁。典型场景是机房退租、上云、自建转云托管,整个集群要换到新环境。这种情况下Topic、消费者组、存量数据全部要移动,而且往往有停机窗口限制,业务不能断太久。这类迁移最复杂,也是本文主要展开的场景。

第二类是集群内节点替换。比如某几台broker要退役、磁盘报警要换盘、或者单纯想扩容横向加节点。这时候集群本身不换,只是把分区从旧节点挪到新节点。很多人不知道的是,这类场景根本不需要跨集群复制,用Kafka自带的kafka-reassign-partitions.sh就能在线完成。

第三类是跨大版本升级。比如从0.10直接升到3.x,除了搬数据,还要考虑协议兼容、消息格式差异、客户端版本是否匹配。这类迁移涉及的东西比较杂,经常需要结合双写或MirrorMaker来做灰度推进。

这三种场景的解决思路完全不同。如果你不先判断类型,很容易选错方案。我一般会先问三个问题:集群要不要整体换?业务允不允许停机?客户端代码能不能改?答案一旦清楚,方案基本就浮出水面了。

1.2 先盘点存量:Topic、分区、副本、消息量、保留策略

Kafka迁移最忌讳的就是“对着一个Topic就开始搬”。动手之前必须盘家底,而且是要盘得非常细。否则你连要准备多少磁盘、复制链路要多宽、数据校验要跑多久都说不出来。

需要盘的东西包括这几类:所有Topic的列表和配置,直接用bin/kafka-topics.sh --bootstrap-server old:9092 --list和--describe就能拿到。注意看每个Topic的分区数、副本因子、retention.ms、cleanup.policy。然后是每个Topic的流量特征,峰值写入速度、单条消息大小、每秒消息数。这个数据最好从broker的JMX监控里拉,如果之前没监控,就得趁迁移前补上,否则后面做容量规划就是拍脑袋。

消费组也要拉一份清单:bin/kafka-consumer-groups.sh --bootstrap-server old:9092 --describe --all-groups,确认有哪些group、分别消费哪些Topic、当前Lag是多少。最后是broker资源情况,磁盘使用率、IO util、网卡带宽,这些决定了你新集群的硬件选型。

我通常会做一个容量计算示例,比如某个订单Topic峰值写入2MB/s,单条消息平均4KB,保留7天,副本因子2。单日数据量是2MB/s乘以86400秒,约172.8GB,7天就是1.2TB,再乘以副本2,得到2.4TB,最后留出30%缓冲,约3.1TB。这个数字才是你给新集群配磁盘的依据。

这里想强调一个经验:Kafka迁移里总数据量其实不是主导因素,真正决定迁移难度的是峰值写入速度和保留时长。这就像搬家,你真正要关心的是每天产生多少新东西、垃圾桶多久清一次,而不是数家里一共有多少件旧家具。

1.3 迁移目标选型:版本、集群规模、磁盘与网络规划

存量盘清楚后,就是给新集群做选型。版本方面,现在新项目基本都选3.x以上,并且建议直接用KRaft模式,也就是不依赖Zookeeper的那种部署方式,运维负担明显小很多。如果团队对KRaft还不太信任,3.x也支持传统ZK模式,可以平滑过渡。关键不是追求最新,而是选一个你们团队真正能运维得起来的版本。

集群规模要根据分区总数和峰值吞吐来估算。我的经验值是单broker分区数控制在2000以内比较稳,超过5000就会有明显的性能风险。副本因子生产环境至少2,推荐3。磁盘方面,Kafka是顺序写盘,机械硬盘也能扛一定吞吐,但延迟敏感的业务还是优先SSD,尤其是要做消息检索或重放的场景,SSD提升非常明显。网卡至少千兆起步,跨机房对等连接建议专线或万兆,否则同步延迟会很难看。

新集群的Topic规划有个原则:迁移阶段尽量和旧集群保持一致,分区数、副本数、retention这些先照搬,让客户端和业务逻辑零改动。等迁移完成、系统稳定运行一段时间后,再根据实际流量做分区调整或参数调优。一上来就“顺手优化”分区设计,往往会引入额外变量,出问题了很难定位是新环境的问题还是改造的问题。

2. 迁移方案对比:双写、MirrorMaker 2、还是分区重分配

2.1 应用侧双写:最土但最可控

应用侧双写是最“不优雅”但最可控的迁移方式。思路很简单:改造生产者的发送逻辑,让每条消息同时写入新旧两个集群。业务通过配置中心控制写入开关,按Topic或者按实例灰度切换,整个过程完全由你自己掌控节奏。

双写的优点很明显:无需额外组件,不需要部署MirrorMaker,数据同步的时效性完全取决于业务发送链路,理论上比任何异步复制都快。另外切换非常灵活,今天可以只让10%的消息打到新集群,观察稳定后再加大比例。

但缺点也必须提前想清楚。第一,业务代码一定要改,如果系统里几十个服务都在直连Kafka,改造成本会很高。第二,消息重复不可避免。两条链路都可能出现部分成功的情况,比如新集群写失败但旧集群写成功了,下游消费时如果不做幂等,必然会有重复消息。第三,双倍流量会对生产链路和下游造成额外压力,一些对延迟敏感的核心链路需要评估是否承受得住。

实际操作时,我习惯在Producer外面包一层路由,底层还是原生的KafkaProducer,只是发送前根据路由规则复制一份到另一个Producer实例。这里有个很容易踩的坑:两个集群都写失败时怎么办?降级策略一定要提前定好。我的经验是,新集群写失败只记录日志和指标,不影响主链路;旧集群写失败则直接抛异常,触发原有告警,因为那才是当前真正在用的链路。

双写方案适合业务可改造、迁移时间充裕、希望全程能灰度控制的团队。它不适合那种几十个老系统都直连Kafka、代码已经没人敢动的场景,那种场景还是老老实实用MirrorMaker。

2.2 MirrorMaker 2:跨集群复制的官方方案

如果你不想改业务代码,那MirrorMaker 2(简称MM2)基本是首选。它是Kafka官方自带的跨集群复制工具,从2.4版本开始就是标准答案。它的本质是一个跑在Kafka Connect框架上的应用,用内置的Source Connector消费源集群的Topic,再通过内置Producer把数据写入目标集群。

MM2比老版MirrorMaker强的地方在于,它不只是复制数据,还会同步消费组offset、Topic配置、ACL等元数据。特别是checkpoint机制,它会定期把源集群消费组的位点记录到目标集群的专用Topic里。这个能力对迁移来说非常关键,它决定了你切换消费者后,能不能从旧集群的断点继续消费,而不是从头消费或者从末尾消费。

用MM2做迁移的最大优势就是业务零改造,不用碰任何Producer和Consumer代码,只要部署一个MM2进程就能把整条数据管道搭起来。同时它支持双向复制,也符合容灾演练的场景。缺点方面,MM2是异步复制,目标集群数据始终存在一定延迟,几秒到几十秒都有可能,具体取决于网络和负载。此外复制过程中也可能产生重复消息,下游消费最好有幂等兜底。

如果只是迁移需求,我通常只开单向往目标集群复制,不开反向。双向复制在容灾场景才有意义,迁移阶段开双向容易造成消息循环,给自己增加不必要的复杂度。

2.3 分区重分配:集群内节点迁移的正确打开方式

这里单独说一下集群内节点替换的场景。很多人一听到迁移就想到跨集群复制,其实如果只是要换掉某几个broker,用Kafka原生的分区重分配工具就够了,完全不需要MirrorMaker。

原理上,kafka-reassign-partitions.sh会把指定分区的数据在broker之间做增量复制,等新副本追平到leader之后,切换leader,再删除旧副本。整个过程是在线完成的,不需要停机,但会在迁移期间产生额外的磁盘和网络IO。所以我的建议是尽量在业务低峰期执行,并且用--throttle参数限速,比如限到50MB/s,避免把集群IO打满。

操作流程一般是三步。第一步用kafka-reassign-partitions.sh --generate生成候选迁移方案,它会根据当前分配情况输出一个JSON分配方案;第二步把方案改到只包含你想迁移的Topic和分区,然后用--execute执行;第三步用--verify观察迁移状态,直到所有分区都显示completed。

这里有个重要经验:重分配一个批次不要挪太多分区,一批一批来。你一次把所有分区都搬过去,中间某个节点挂了,整个集群的稳定性都会受影响。我习惯一次最多挪几十个分区,跑完一批确认稳定再做下一批。整个过程虽然慢,但稳。

2.4 如何选:停机窗口、一致性、改造成本

方案选型的核心就三个变量:停机窗口、数据一致性要求、业务改造成本。我整理了一个简单的对比表:

方案停机时间业务侵入数据一致性复杂度适用场景
冷拷贝停写需要无高低允许停机的小集群
应用双写不需要高中(可能重复)中业务可改造、可灰度
MirrorMaker 2不需要低中(异步)中高整体跨集群迁移
分区重分配不需要无高低集群内节点替换

如果停机窗口比较富裕,其实可以选最土的冷拷贝方案:停掉所有生产写入,把Kafka数据目录或者用工具搬迁存量数据,然后在新集群启动消费。这个方案数据一致性最高,操作也最简单,问题就在于业务要能接受停机时间,哪怕只有十几分钟,也不是所有系统都扛得住。

生产环境多数不能接受长时间停机,所以“MirrorMaker 2/双写 + 消费组offset同步”是主流组合。面试时如果被问到Kafka怎么做迁移,先问清楚能不能停机、要不要改代码、是跨集群还是集群内换节点,然后再给方案,基本就能加分。高手不是会敲一条命令,而是会做方案取舍。

3. 实操:用 MirrorMaker 2 做一次跨集群迁移

3.1 环境准备:本地 Windows 快速搭一套测试集群

很多同学本地是Windows,想先搭一套环境验证方案。Kafka官方其实没有提供Windows安装包,但二进制包解压后是可以直接跑的,只要把JDK配好。Kafka 3.x以上推荐用KRaft模式,不用装Zookeeper,比很多老教程舒服得多。

步骤不复杂。先下载Kafka二进制包,比如kafka_2.13-3.6.2,解压到C:\kafka。然后把JAVA_HOME指向JDK 11或17,PATH里加好bin目录。接着编辑config\kraft\server.properties,关键配置改成这样:

process.roles=broker,controller node.id=1 listeners=PLAINTEXT://localhost:9092,CONTROLLER://localhost:9093 advertised.listeners=PLAINTEXT://localhost:9092 controller.quorum.voters=1@localhost:9093

然后打开命令行,先执行bin\windows\kafka-storage.bat random-uuid生成一个集群ID,保存下来。再执行bin\windows\kafka-storage.bat format -t <uuid> -c config\kraft\server.properties格式化存储目录,最后用bin\windows\kafka-server-start.bat config\kraft\server.properties就能把单节点Kafka跑起来。

想模拟新旧两个集群,就再复制一份目录,把server.properties里的node.id改成2,端口改成9192和9193,再格式化启动一次。两个本地集群就能跑通后面的MM2流程。Windows下有几个常见坑:路径带中文会报错,JDK版本太老起不来,改了listeners但忘了改advertised.listeners会连接失败。这些细节排查起来很耗时间。

这里的重点是搭一个能复现迁移流程的测试环境,不是为了性能验证。生产环境还是建议Linux裸机或容器化部署,Windows当个验证工具足够了。

3.2 MirrorMaker 2 的配置与启动步骤

本地假设有两个集群:source监听localhost:9092,target监听localhost:9192。我们需要把source的Topic复制到target。

新建一个mm2.properties配置文件,核心内容如下:

clusters = source, target source.bootstrap.servers = localhost:9092 target.bootstrap.servers = localhost:9192 # 单向复制:source -> target source->target.enabled = true source->target.topics = .* target->source.enabled = false # 自动创建目标集群Topic source->target.topic.auto.create = true # 同步消费组offset source->target.emit.checkpoints.enabled = true source->target.sync.group.offsets.enabled = true

启动命令是bin/kafka-mirror-maker2.sh config/mm2.properties,Windows下就是bin\windows\kafka-mirror-maker2.bat。启动后到target集群里看,会发现多了一批新Topic,名字默认带source.前缀,比如source.orders。

这个前缀是MM2故意设计的,目的是防止双向复制时消息循环。但对迁移场景来说,我们通常希望目标集群的Topic名和源集群完全一致,这样客户端切换配置时最省事。要实现这一点,需要自定义一个ReplicationPolicy。

写一个Java类,继承DefaultReplicationPolicy,重写formatRemoteTopic方法:

import org.apache.kafka.connect.mirror.DefaultReplicationPolicy; public class NoPrefixReplicationPolicy extends DefaultReplicationPolicy { @Override public String formatRemoteTopic(String topic) { return topic; } }

编译成jar,放到Kafka安装目录的libs下,然后在mm2.properties里指定replication.policy.class=NoPrefixReplicationPolicy。这样复制出来的Topic名就和源集群一致了。

这里必须提醒一句:一旦用了无前缀策略,千万别再开双向复制,否则两个集群互相复制,消息会无穷循环。单向复制不受影响。另外,Topic比较多时,建议先用正则限定范围,比如source->target.topics = orders\|payment\|user_events,验证没问题后再放开全量,避免一开始就复制一些不需要的Topic。

3.3 验证数据完整性:消息数、位点、时间戳

复制跑起来之后,最容易被忽视的就是验证环节。我见过有人启动MM2后看到两边Topic都存在就宣布迁移成功,结果过几天业务方发现有消息对不上,因为某些分区offset没追上。

第一层验证是对比各分区offset。在source集群执行:

bin/kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list localhost:9092 --topic orders --time -1

target集群也执行相同命令,然后把每个分区的log end offset逐行对比。更稳一点的做法是写个脚本,先拿每个分区的log start offset,再拿log end offset,两者之差就是该分区的消息总量,然后两端对比。注意不要只看总和,分区粒度上的对比才能暴露具体是哪个分区出了问题。

第二层验证是抽样比对消息内容。用console consumer把部分分区的消息导出,或者写一个简单的Kafka Consumer,对每条消息的key和value做哈希,放到一个Map里,再对另一端做同样操作,比对是否有差异。抽样不需要全量,但每个分区都要覆盖到。这一步能发现offset对得上但内容不对的场景,比如序列化问题或者复制链路中间处理逻辑有问题。

第三层验证是消费端位点。在source集群查看业务消费组的当前位点:

bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-consumer-group

如果MM2的checkpoint同步开启,target集群会有一个同名group.id的消费组,位点应该接近source。这个值决定了切换消费者后,能不能从旧集群断点继续消费。

实操心得:MM2是异步复制,目标集群总会比源集群慢几百毫秒甚至几秒,这是正常现象。验收标准不要写成“两边完全一致”,而是定义“差距在可接受范围内且不持续增大”,比如5秒以内就算通过。

3.4 流量切换与客户端割接

数据校验通过后,进入切换环节。这里的顺序非常重要,我强烈建议:先切消费者,后切生产者。

为什么不能先切生产者?因为如果先把生产者切到新集群,旧集群的消费者还在旧集群消费,新写入目标集群的消息旧消费者完全看不到,业务直接就断了。反过来先切消费者,因为MM2持续把源集群新消息同步到目标集群,消费者切过去后虽然可能短暂滞后,但很快会追平。

具体操作分四步。第一步,确认target集群已有同名消费组,并且offset与source集群接近,这样才能续接消费位点。第二步,把消费者客户端的bootstrap.servers从旧地址改到新地址,分批灰度切换,改完后立刻观察新集群这个消费组的lag,确认是从同步的位点继续消费而不是从头或从尾部开始。第三步,消费者全部切完并稳定后,再切生产者,同样分批灰度。第四步,全部切完后,MM2可以继续运行一段时间作为备份链路,但要注意复制方向不能开反向,否则会产生循环复制。

这里有个很常见的坑:消费者切过去后,发现位点不对,从开头开始消费了。原因一般是sync.group.offsets.enabled没开,或者源目标和客户端里的group.id不一致,或者目标集群这个group之前被手动消费过导致checkpoint失效。检查这三个点基本能定位。

切换过程中如果用了自定义ReplicationPolicy去掉了Topic前缀,还要注意一点:客户端不能同时连接两个集群去订阅同一个名字的Topic,元数据会互相覆盖,导致路由混乱。所以切换期间最好是同一进程内只指向一个集群,不要新旧集群地址同时挂在同一个客户端配置里。

4. 迁移中的兼容性与坑:大消息、延迟、版本差异

4.1 单条 1MB 大消息的配置与验证

搜索Kafka相关问题时经常能看到类似“kafka 接收1m”的关键词,这说的就是Kafka默认单条消息大小上限是1MB。很多业务在迁移阶段才暴露出这个问题:旧集群可能早就被运维调大了参数,新集群还是默认值,一切过去立刻报RecordTooLargeException。

需要对齐的参数有好几层,我先列个表:

位置参数默认值调整方向
brokermessage.max.bytes1048576调到目标单条上限,如10485760
brokerreplica.fetch.max.bytes1048576必须大于等于message.max.bytes
brokersocket.request.max.bytes104857600按需调大
producermax.request.size1048576大于等于单条消息上限
consumermax.partition.fetch.bytes1048576大于等于单条消息大小
consumerfetch.max.bytes52428800按批量消费需求调整

我经常看到有人只改了broker的message.max.bytes,忘了改replica.fetch.max.bytes,结果副本同步时一直报错,副本长期处于UnderReplicated状态。这两个参数必须一起调整。

验证方法很简单,写一段Producer代码,发送一条2MB的字节数组到目标集群,能正常发送、能被消费,参数才算配好:

Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9192"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer"); props.put("max.request.size", "10485760"); Producer<String, byte[]> producer = new KafkaProducer<>(props); byte[] bigPayload = new byte[2 * 1024 * 1024]; producer.send(new ProducerRecord<>("orders", "migration-key", bigPayload)).get(); producer.close();

迁移前把新旧集群这套参数全部对齐并实际用大消息测一遍,能节省非常多线上排障时间。图片、日志大字段、埋点大JSON在真实业务里非常常见,不是极端场景。

4.2 迁移后消息延迟变高的排查思路

切到新集群后,业务反馈“消息延迟高”是高频问题。出现这种情况不要急着怪新集群性能差,先按环节拆解:到底是生产者到broker慢,还是broker到消费者慢。

我排查时会先看broker的CPU、磁盘IO、GC。磁盘IO是很多迁移翻车的地方,旧集群用了SSD,新集群却配了普通HDD,顺序写差距可以达到好几倍。然后是生产者侧指标,batch size、linger.ms、compression.type、acks。acks=all时,如果副本数从2变成了3,确认时间自然会变长;linger.ms=0时,小消息每一条都单独发送,网络往返成本非常高。

消费者侧要关注max.poll.records、max.poll.interval.ms、session.timeout.ms。消费者处理不过来,lag就会涨,整体表现就是消息延迟高。

还有一个特别隐蔽的坑:跨机房网络。新旧集群不在同一个机房时,客户端切换后的RTT会明显上升,尤其是同步发送加acks=all,延迟直接翻倍。物理距离带来的延迟不是调参能完全消除的,尽量让新旧集群在迁移窗口内处于同一网络区域,或者至少用专线打通。

我整理了一个快速排查表:

现象可能原因查找方式对策
producer send耗时高linger.ms太小/batch太小producer metric batch-size调整batch.size和linger.ms,考虑开启压缩
消费lag持续涨消费者处理慢/分区不够kafka-consumer-groups --describe增加消费者实例或分区,优化消费逻辑
副本同步异常replica.fetch参数没对齐kafka-topics --describe 看ISR对齐大消息相关参数
端到端延迟大网络RTT高/跨机房ping、专线路由监控迁移期间尽量同机房,或调整生产者确认模型

另外还有一个现象:客户端改了bootstrap.servers指向新集群后,可能还连着一部分旧broker。这是因为客户端会缓存元数据,不会立刻丢弃旧的broker连接。出现这种情况时,最直接的办法是重启消费者进程,等metadata刷新后再观察。

4.3 客户端版本兼容与序列化

迁移前还要盘一下所有客户端的Kafka版本。Kafka服务端对老客户端有兼容性策略,但也不是无底线地兼容。0.10.2以上的客户端连3.x基本没问题,0.9或者更早的版本,新集群默认可能会拒绝老协议。

除了协议版本,消息格式也是一个隐蔽问题。新集群默认消息格式是v2,如果你的消费者仍然是老版本kafka-client并且没有做兼容设置,可能出现反序列化或CRC校验问题。稳妥做法是迁移前把客户端统一升级到0.11以上;如果确实无法升级,可以在新集群设置log.message.format.version为旧版本对应的值,但这个参数在新版本中有废弃趋势,终归要升级客户端。

序列化方面容易被忽略。Kafka本身不关心消息内容格式,但迁移验证时如果你按JSON或AVRO去解析,就必须确保新旧两边的schema一致。特别是用了Confluent Schema Registry的场景,新集群客户端需要指向新的registry地址,同时把schema同步过去,否则消费者反序列化直接报错。

我在实际迁移中一定会做一张“客户端清单”,列清楚每个服务用的kafka-clients版本、序列化方式、是否使用Schema Registry、消费组ID是什么。别嫌这张表麻烦,等出了问题再去翻代码,成本翻十倍都不止。

5. 迁移后的观测与运维:UI界面、监控、性能对比

5.1 Kafka UI 工具推荐

顺带回答那个高频问题:Kafka有没有UI界面?答案是不仅有,而且选择非常多。迁移验证阶段我肯定会装一个UI,直观看到新旧两个集群的Topic、分区和消息情况,比敲命令行舒服很多。

常见选择:

工具类型特点适用场景
Kafka UIWeb多集群、消息查看、consumer管理,社区活跃日常运维首选
Offset Explorer桌面客户端查看Topic、分区、offset非常方便本地快速调试
KafdropWeb轻量,看消息直观临时排查
CMAK (Kafka Manager)Web老牌管理工具,偏集群管理老版本集群迁移期
Prometheus + GrafanaWeb监控指标监控必备生产环境告警

经验提醒:生产环境的UI工具建议只做只读用途,不要在UI里直接改Topic配置、删消费组。这些操作很容易误点,尤其Kafka UI新版本交互改动频繁,我见过同事点错按钮把消费组重置了的情况。要改配置还是用命令行加权限控制比较稳妥。

5.2 监控指标与性能对比

迁移完成后,不是看一眼lag能消费就算完。正确姿势是建立一套核心指标监控,把新旧集群迁移前后做个对比,确认新集群确实承接住了原流量。

重点盯这几个指标:BytesIn/BytesOut表示集群吞吐,确认有没有流量异常;MessagesInPerSec是每秒消息数,和业务预期对照;RequestHandlerAvgIdlePercent是请求线程空闲率,长期低于30%说明broker线程快被耗尽,需要调大num.io.threads;NetworkProcessorAvgIdlePercent同理。UnderReplicatedPartitions是副本同步落后分区数,长时间大于0说明某台broker写入慢或网络有问题;OfflinePartitionsCount必须是0;IsrShrinksPerSec和IsrExpandsPerSec如果频繁跳动,说明broker性能波动大。

用Kafka exporter配合Grafana模板就能搭一套,不用自己从零画图。迁移后第一天要重点盯UnderReplicatedPartitions和IsrShrinksPerSec,因为新集群的副本同步往往是最弱的环节。

性能对比也可以用压测脚本验证。Kafka自带工具:

bin/kafka-producer-perf-test.sh --topic orders --num-records 1000000 --record-size 1024 --throughput 10000 --producer-props bootstrap.servers=localhost:9192 linger.ms=10 batch.size=65536
bin/kafka-consumer-perf-test.sh --bootstrap-server localhost:9192 --topic orders --messages 1000000 --reporting-interval 1000

压测前先确认Topic分区数足够,如果分区数太少,吞吐会被单个分区上限卡死。压测结果不一定代表生产真实水平,但用来做新旧集群的横向对比很有参考价值。

5.3 迁移验收与回滚兜底

验收维度我习惯列一个清单,逐项打勾。新旧集群各项Topic配置要一致,包括分区数、副本数、retention、cleanup.policy、min.insync.replicas。新旧集群各分区offset差值要低于设定阈值。业务消费组在新集群运行正常,lag稳定在低位。核心链路场景要人工验证通过,比如下单、支付回调、报表任务这些真实业务动作,不能只看监控。监控告警全部切换到新集群,旧集群告警可以保留但降噪。

再说回滚。我的原则是:旧集群不要急着下线。最稳的回滚策略是保留旧集群和复制链路至少一个完整保留周期,假设retention设置的是7天,就留满7天。如果切换后发现问题,生产者切回旧集群,消费者也切回旧集群,位点还在,业务能迅速恢复。新集群这段时间积累的数据可以手工补数,尽量减少损失。

回滚期间最忌讳的操作是删Topic、删消费组、清空offset。这些动作一旦做了,旧集群就真的回不去了。所以下线旧集群前,一定要让负责人确认:数据已确认完好,业务已稳定运行超过一个保留周期,再动手清理。

还有一个小细节:迁移结束时,MM2进程不要用kill -9粗暴停掉。如果你还在用默认复制策略,它可能正在运行checkpoint同步,直接杀掉容易留下未完成的内部状态Topic。正确做法是先把mm2.properties里的复制开关改成false,或者停掉connect任务,等内部状态Topic清理完再退出进程,避免下次启动时出现offset错位。

最后,数据校验脚本千万别删,留在团队里。每次迁移或者大版本升级都可能用得上。我自己的体会是,Kafka迁移真的不是看谁命令敲得熟,而是看谁能把停机窗口、数据校验、回滚预案这些“脏活”想在前面。把这套思路掌握住,迁移过程再复杂,心里也不会慌。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询