上周刚帮朋友公司做完一次Kafka迁移,从两套Kafka 2.8集群跨机房搬迁合并成一套新的Kafka 3.2集群。接到这个需求的时候,我其实也犯过和大多数人一样的懒:觉得Kafka迁移就是把数据拷贝过去,然后让客户端改个连接地址就行。真正动手做下去才发现,一个看似简单的“搬迁”背后,牵扯的是主题分区对齐、消费位点映射、消息延迟控制、认证策略统一、镜像同步压力等一系列一连串的问题。整个过程前后折腾了七天,中间踩了不少坑,也总结了一套能复用、能回滚、能验证的完整思路。这篇文章就把这次Kafka迁移从方案选型、前期摸底、执行切换、问题排查到验收运维的全过程写一遍,给正准备做Kafka集群替换、机房搬迁或跨环境同步的人当个参考。
先说一个共识:Kafka迁移最怕的不是技术方案不够炫,而是三个核心问题没解决——消息不能丢、消费不能乱、切换不能抖。任何迁移方案本质上都是围绕这三个目标做取舍。方案没有绝对的好坏,只有适不适合你的业务体量和团队协作方式。我结合这次实际项目的情况,把三种主流方案都分析一遍,你看完就知道自己该走哪条路。
1. 迁移前先想清楚:三种主流方案怎么选?
1.1 停机冷迁移:适合小规模场景的兜底方案
所谓停机冷迁移,就是在凌晨业务低峰期申请一个停机窗口,把所有的生产者和消费者全部停掉,然后直接把旧集群的日志分段文件拷贝到新集群的对应目录,或者借助kafka-reassign-partitions这类工具把分区数据搬到新集群,最后再启动客户端修改连接地址。
这个方案的逻辑最简单,因为生产端和消费端都停掉了,不存在两边集群同时写入导致的乱序、重复问题。但它的缺点也非常致命:对线上大集群基本不可行。我之前遇到过数据量上百GB甚至上TB的场景,光拷贝日志文件再加载、校验副本,一个晚上根本不够,而且拷贝期间但凡有任何一条消息只写进了旧集群没同步过去,迁移后数据就是不完整的。你很难在狭小的停机窗口内完成数据量核对、条数校验、消费位点确认这一整套动作。
我的判断是:除非你只有测试环境级别的数据,或者业务本身允许凌晨停服几个小时,否则不要把冷迁移作为线上主力方案。它可以作为数据量极小场景下的兜底手段,但当一个迁移项目涉及几十个topic、几百个分区的时候,冷迁移等于给自己挖坑。
1.2 协议层镜像同步:MirrorMaker的适用边界
镜像同步的思路是把新集群先部署好,然后通过MirrorMaker 2或者Confluent Replicator这类工具,将旧集群里的topic实时同步到新集群。消费者切换过去之前,新集群里已经囤好了一份完整的数据,切换动作变成一个普通消费组改地址的操作。
这个方案最大的优势是不需要改业务代码,对topic多、分区多的大集群尤其友好。我这次迁移的旧集群有接近300个topic,每天数据量几百GB,如果走双写方案要改几十个服务,根本不是一两天能推进的事情。镜像同步把迁移变成了一个“后台数据搬运任务”,业务方只需要等我们把数据对齐,然后配合切换即可。
但镜像同步也有代价:数据同步有延迟,正常情况几秒到几十秒;如果业务对端到端顺序有强要求,镜像链路的调优会非常讲究;而且MirrorMaker 2默认的复制策略会给topic自动加集群别名前缀,比如旧集群别名是old,同步过去的topic会叫old.topic-name,如果不做配置,消费者切到新集群后根本找不到原来的topic名称。
1.3 双写方案:最稳妥但代价是代码改造
双写是指业务在发送消息时同时投递到新旧两个集群,消费者先切到新集群消费,跑几天确认稳定后,再让生产者逐渐停掉旧集群的写入,最后下线旧集群。这个方案的切换粒度最细,可以一个消费组一个消费组地灰度切换,出了问题时回滚路径也最清晰。
它的缺点也很明显:需要业务方配合改代码,而且双写期间要处理重复消息、跨集群幂等、延迟对比等一系列问题。如果你的Kafka是公共基础设施,多个业务团队共用,想同时说服所有业务方配合改造,通常很难推进。双写适合那些数据链路完全在自己可控范围内的核心服务,比如你就是某个订单系统的负责人,服务的生产消费逻辑都能自己改,那用双写确实是最稳的路线。
三种方案放在一起对比的话:
| 迁移方案 | 业务代码改造 | 数据实时性 | 回滚难度 | 适用场景 |
|---|---|---|---|---|
| 停机冷迁移 | 无 | 迁移期间完全停服 | 回滚复杂 | 小数据量、可接受停服 |
| 镜像同步 | 无需改代码 | 秒级延迟 | 回滚简单,保留旧集群即可 | 大规模共享集群、跨机房搬迁 |
| 双写 | 需要改造 | 实时双写 | 回滚最灵活 | 核心业务自控链路、灰度切换 |
我个人在这次的迁移里选择了镜像同步方案,因为团队的场景决定了这就是最优解:topic数量大、业务方多、无法统一改代码,同时可以保留旧集群较长时间用于回滚。
2. 迁移前夜:集群参数盘点与新集群规划
2.1 先给集群做“健康体检”
不管选哪种方案,迁移前要做的事情都一样:先摸清家底。这一步我建议大家不要偷懒,直接上服务器把Kafka的存量数据完整摸一遍。我会做下面几件基础动作:
- 用
kafka-topics --bootstrap-server old-cluster:9092 --list把全部topic列出来,重点记录每个topic的分区数、副本因子、保留策略,以及单topic每天的数据增长量。 - 用
kafka-consumer-groups --bootstrap-server old-cluster:9092 --describe查看所有消费组当前的lag情况,记下迁移前每个消费组应该消费到的位点。 - 记录broker侧的关键配置,比如
log.retention.hours、log.segment.bytes、message.max.bytes、replica.fetch.max.bytes。 - 确认客户端的连接方式,是Plaintext、SASL_PLAINTEXT还是SSL,有没有使用Schema Registry,有没有自定义拦截器。
这些信息直接决定新集群该怎么建。我这次摸底时发现一个很容易被忽略的细节:有一个topic的单条消息体积接近1MB,而新集群默认message.max.bytes只有1MB,换算下来那条消息加上协议开销就会超限。如果不提前发现,迁移后这个topic在业务高峰会持续报“RecordTooLargeException”,所有写入全部失败。我后来把全表topic按最大消息大小拉出来盘点,才发现有两个topic需要单独调大broker侧和topic侧的参数。这就是“先摸底”环节不能省的原因。
2.2 新集群规格与参数规划
摸完家底之后,就要规划新集群的规模了。这里给一个我自己常用的估算方式。
集群总存储 = 所有topic日增长量之和 × 保留天数 × 副本因子。比如假设有10个topic,每个日增长100GB,保留7天,副本因子3,那总存储就是100GB×10×7×3 = 21000GB,也就是大约21TB。考虑磁盘水位不能超70%,建议规划容量至少要30TB以上。如果单台broker挂4块4TB盘、净容量16TB,那么两个broker就已经足够,但还要同时考虑峰值流量和延迟要求,所以实际规划时我会在这个基础上再加一台做冗余。
带宽估算同理:如果迁移期间新集群既要承担镜像同步流量,又要承担业务本身的读写流量,那么网络规划要按照日常峰值的两倍来估算。我的经验是broker的磁盘利用率如果超过90%,副本同步会因为IO抖动出现频繁的ISR收缩和扩张,这时候整个集群的延迟表现都会变差,业务侧会看到明显的消费滞后。
新集群的分区参数也要提前统一。num.partitions、default.replication.factor这些默认值如果两套集群不一样,迁移时创建topic的分区数就对不上,后面会引发按key路由错乱的问题。所以迁移前最好把两套集群的默认参数对比一下,该改的先改。
2.3 版本兼容与认证策略统一
Kafka的跨版本迁移是最容易被低估的风险点。从Kafka 2.x到3.x中间发生了很多变化,比如ZooKeeper依赖在3.x逐步被KRaft替代,很多老的客户端协议也发生了变化。迁移时尽量选择新旧版本差距小的组合,最好不要跨两个大版本以上。这次我是从2.8迁到3.2,中间虽然跨了版本,但因为客户端的版本和协议协商还能兼容,所以推进得还算顺利。如果客户端版本特别老,比如还是0.10或0.11时代的东西,新集群即使能配出兼容模式,也只是临时方案,最终还是得推动客户端升级。
认证策略方面也一样。如果旧集群用的是SASL/PLAIN,新集群就不要图省事改成SASL/SCRAM,除非你提前做好了客户端的改造计划。我见过某个项目切换后发现客户端A还在用老配置连接新集群,结果全部认证失败,业务瞬间变成只读的事故。另外,Kafka的ACL是集群级别的,MirrorMaker不会自动从旧集群把ACL搬过去,需要提前导出再重建,否则下游客户端能连通但没有任何读写权限,这种问题在切换后才会爆出来,排查起来特别被动。
3. 迁移执行:从部署到切换的全流程实录
3.1 新集群部署与验证
新集群部署本身没什么玄学,但要踩准几个点。数据盘要单独挂载,记得调大文件描述符和最大线程数限制;堆内存根据broker所在节点的总内存来配,我一般给broker堆内存设为4到6GB,配合页缓存一起用;操作系统的vm.max_map_count如果太小,运行过程中会出现内存映射不足的异常。部署完成后,先不要急着同步业务数据,我的习惯是新建一个临时topic,生产一批带标记的消息,再消费出来校验一遍,确认broker之间的副本同步正常、没有UnderReplicatedPartitions,然后再继续往下走。
部署完后的验证清单大概是这样的:kafka-topics --describe能看到新创建topic的leader和replica分布正常;用生产者生产消息没报错;用消费者消费出来消息内容一致。这些基础验证用临时topic加造数据的方式最直接,不要一上来就拿真实业务topic去压镜像同步。
3.2 开启镜像同步并验证数据一致性
镜像同步我用的是Kafka自带的MirrorMaker 2,配置中心点是把MirrorSourceConnector指向新旧集群,并指定需要同步哪些topic。核心配置大致长这样:
# mm2.properties clusters = old, new old.bootstrap.servers = old-host:9092 new.bootstrap.servers = new-host:9092 old->new.enabled = true old->new.topics = .* # 复制因子与原集群保持一致 replication.factor = 3 # 保持新集群topic名称不追加集群别名前缀 replication.policy.class = org.apache.kafka.connect.mirror.IdentityReplicationPolicy # 关闭ACL自动同步,ACL迁移我们手动做 sync.topic.acls.enabled = false checkpoint.interval.ms = 5000这里需要特别强调两点。第一,MirrorMaker 2默认的复制策略是DefaultReplicationPolicy,会给同步过去的topic加上源集群别名前缀,比如old.order-log,这会让下游消费者无所适从。所以我在这里显式指定了IdentityReplicationPolicy,让topic名称保持不变。第二,MirrorMaker自己运行时也会创建一些内部topic,比如heartbeat和checkpoint,这些内部topic不要一股脑同步到新集群,否则会污染业务数据目录。
配置写好后,启动命令很简单:
bin/connect-mirror-maker.sh config/mm2.properties但不要一启动就同步全部300个topic,我建议先挑几个小topic做验证,确认同步过去的topic分区数、副本数、消息内容都和原集群一致,观察同步延迟在什么量级,然后再把全部topic放开。同步期间还要持续关注MirrorMaker的Consumer Lag,一旦发现同步速度追不上生产速率,就要考虑增加num.streams或者限流参数来调节。
3.3 消费端流量切换与稳定性观察
数据同步到新集群且验证没问题之后,才开始切消费端。切换顺序我强烈建议是“先切消费、再切生产”,而且消费端也要挑一个对延迟不那么敏感的消费组先切过去。切完后观察一段时间,比如半小时到一小时,对比新集群消费到的位点是否和旧集群对齐,有没有消息积压。
切换消费端的本质是改消费组连接的bootstrap.servers。如果新旧集群的数据完全同步,消费组理论上可以从新集群继续消费。但这里有个关键问题:如果用MirrorMaker同步,新集群的消费位点和旧集群并不是同一个位点体系。最简单的处理方式是先记录迁移前各消费组的位点,切换后用kafka-consumer-groups --reset-offsets把新集群的消费位点重置到迁移前的位置。比如:
kafka-consumer-groups.sh --bootstrap-server new-cluster:9092 \ --group order-consumer \ --reset-offsets --to-offset <迁移前记录的位点> \ --execute消费端切换稳定之后,再分批切换生产者。生产者切换建议按业务线分批次推进,不要在一个时间点把全部生产流量都压到新集群。每次切换后盯两条指标:消息积压量(lag)是否持续走低,消费速率是否和切换前一致。只要这两条正常,基本就可以推进下一步。
整个灰度切换期间,旧集群一定要保持运行状态。如果新集群出现处理问题、数据同步跟不上或者业务指标异常,只需要把连接地址改回旧集群,就能快速回滚。我见过比较倒霉的案例,迁移后第三天发现了消费延迟数据依赖旧集群,但旧集群已经被下架了一半,最后只能找备份恢复,代价非常大。
4. 迁移中的那些坑:常见问题与排查思路
4.1 主题分区数不一致导致的数据错位
镜像同步完成后,我检查新集群的时候发现一个topic的分区数是16,而旧集群明明是24。原因是有个运维同事提前在新集群手工创建了这个topic,默认使用了新集群的num.partitions=16。表面上看消息数量和消息内容都没少,但下游按key聚合统计时结果全部对不上。
这是因为Kafka的消息路由是key的哈希对分区数取模。分区数变了,同一个key的消息就会落到不同的分区,如果有下游逻辑按partition维度处理数据或者做局部排序,那结果一定是乱的。解决方案很直接:删除新集群里预先创建的错误topic,让MirrorMaker按源端拓扑自动重新创建;或者用kafka-topics --alter先把分区数改成一致。这里我建议优先删除重建,因为--alter改分区数虽然可行,但会触发数据重分布,在同步链路里容易被忽视频繁的副本迁移对性能的影响。
4.2 消费位点与消息延迟同时报警
消费组切换到新集群后,最吓人的现象就是监控面板上的lag一直在涨。我第一次遇到这种情况的第一反应是新集群处理能力不行,后来仔细排查才发现问题出在消费参数上。
当时切过去之后,lag从静止状态一路涨到十几万条,查看消费者日志每批拉取的数据量小得可怜。后来发现新集群消费者组用的max.poll.records太小,而且每条消息处理过程中还做了一次远程网络调用,处理速率远低于生产速率。把max.poll.records调大,给session和心跳参数留足余量之后,lag才开始慢慢消化。所以遇到延迟高时,不要急着怪集群,先分两步排查:第一步看lag是整体都高还是个别分区高。整体都高大概率是消费能力不足,个别分区高大概率是热点分区分布不均衡。第二步看消费者实例数和单条消息处理耗时,分区数不够的时候加消费者实例其实没用。
还有一种情况特别容易误导人:MirrorMaker同步期间,如果生产速度短时间暴涨,新集群看到的消息是滞后于旧集群的,这时候它本身就会显示消费lag很高,就算消费组处理能力再强,也只能等数据同步追上来。
4.3 跨版本时的协议不兼容
跨版本迁移里最常见的现象是老客户端连上新集群后持续报Unsupported version exception,或者生产者报错找不到topic。原因是Kafka客户端和服务端之间有协议版本协商,跨大版本时如果客户端版本太老,即使能连上也无法使用新特性,甚至直接失败。
应对方式是在迁移前统计全部客户端的版本。如果发现老版本客户端比例很高,优先考虑先升级客户端再迁移,不要把兼容模式当成长期方案。新集群在启动时也可以临时设置message.format.version来兼容旧的消息格式,但这种向后兼容模式会限制新集群使用新特性,后面还是要找时间改回来。
4.4 镜像同步带来的磁盘与网络压力
MirrorMaker跑起来之后,新集群的broker磁盘写入和网络流量会非常夸张,因为在业务写入之外又叠加了一整份镜像同步流量。这个阶段如果生产者的延迟也跟着变大,问题可能不在业务侧,而是镜像同步把新集群的带宽吃满了。
我遇到过镜像任务的producer压缩设置成了none,结果整个机房出口带宽全被占满的情况。把压缩模式改成lz4或zstd之后,带宽立刻降下来,业务延迟也恢复正常。镜像同步建议选在业务低峰期开启,先让数据快速拉齐;之后在业务高峰期保持低速追赶即可。观察新集群broker的网卡监控和磁盘吞吐,就可以找到一个合适的同步速度,没有必要一直满速跑。
4.5 常见问题速查表
| 现象 | 常见原因 | 排查与处理方式 |
|---|---|---|
| 同步后topic数量对不上 | 通配符配置遗漏内部topic | 核对镜像配置,过滤heartbeat、checkpoint |
| topic名称变了 | 默认复制策略加了集群别名 | 配置IdentityReplicationPolicy |
| 切换后消息路由错乱 | 新集群topic分区数与源端不一致 | 删除重建或--alter对齐分区数 |
| 消费组lag持续上涨 | max.poll.records太小、处理速率不足 | 调大拉取参数、扩容消费者实例 |
| 老客户端连接报错 | 客户端版本与broker协议不兼容 | 升级客户端,临时配置兼容格式 |
| 迁移期间业务延迟变高 | 镜像同步占满带宽或磁盘IO | 开启压缩、降低并发、错峰同步 |
5. 迁移后的验收清单与运维要点
5.1 数据完整性校验
迁移完成后,很多人觉得数据同步过去了就算结束,其实验收才是最容易暴露问题的一步。我的验收习惯是做三核对:第一,topic数量一致,用kafka-topics --list对比新旧集群的topic列表;第二,每个topic的分区数和副本数一致,用kafka-topics --describe逐项比对;第三,关键topic抽样对比消息总条数和最后一条消息的时间戳,用kafka-get-offsets查各个分区的最新位点。
除了消息条数,消费组位点的核对比条数更可靠。切完消费组后跑一段时间,确认没有消费到重复数据或者丢数据,这个比单纯看条数更能说明问题。另外还有一个我常用的土办法:迁移前后把各topic的数据目录总大小做个diff,偏差超过一定阈值就回头排查是不是有分区没同步到位。
5.2 监控与日常运维
迁移后的一周内,我重点盯这几个指标:ByteIn、ByteOut、UnderReplicatedPartitions、OfflinePartitions和Consumer Lag。这些指标可以通过JMX暴露给监控系统,也可以用Kafka自带的命令行工具快速查看。
如果不想每次都上服务器敲命令,可以装一个Kafka UI工具,比如Kafka UI、AKHQ或者Kafka Eagle,都能直接在浏览器里看到topic列表、分区分布、消费组lag这些信息。我个人建议UI工具用于日常巡检没问题,但做迁移验证或数据一致性确认时,还是用命令行最保险,因为UI和后台API之间还有一层转换,信息不一定完整。
5.3 旧集群清理与文档沉淀
旧集群不要急着下线,我建议保留至少一周甚至更长。迁移后跑几天没出问题,再开始走下线流程。下线前,把旧集群的topic清单、消费组、配置备份导出一份,存到文档里。
整个迁移过程也要整理成一份可回放的操作手册:每一步谁操作、什么时间、有哪些观察指标、出了什么问题、怎么解决的。这一步看起来啰嗦,但对后续团队接手非常关键。至少我这次迁移做完之后,把操作手册丢给另一个同事,他做第二次迁移的时候基本没怎么踩我踩过的坑,这就是文档的复利。
最后说一点个人体会。做Kafka迁移,最怕的不是技术方案选错,而是没有按阶段验证就急着切流量。我经历过一次晚上十二点切完所有消费端,凌晨两点全链路报警,才发现消费位点压根没对上,最后回滚到旧集群重新拉数据,折腾到天亮。从那以后我给自己定了个规矩:无论多急,每一次切换之前都要先把“数据一致性验证”这一步走完,宁可在旧集群上多跑一天,也不在切换后发现丢数据。Kafka迁移的核心其实不是“迁移”,而是“验证”。你做完一次这样的迁移,才能真正理解这句话。