ARTICLE DETAIL

资讯详情

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

Kafka集群迁移实战:基于MirrorMaker2的平滑切换方案

Kafka集群迁移实战:基于MirrorMaker2的平滑切换方案 在Kafka运维和架构改造的日常里集群迁移算不上高频但每次遇到都是硬仗。尤其是跨机房搬迁、版本升级、或者干脆把自建集群整体挪到云上稍有不慎就是丢消息、重复消费、业务链路中断的连环事故。我经手过好几次Kafka集群迁移从最早用Kafka内置的kafka-reassign-partitions配合停机窗口硬切到后来用MirrorMaker2做平滑迁移踩过的坑和总结出的方法论都不少。今天这篇重点说MirrorMaker2以下简称MM2的集群迁移实施方案。它本身是Apache Kafka 2.4.0起随Kafka发行版一起发布的跨集群复制工具比老版MirrorMaker强在支持topic自动同步、消费组偏移同步、分区镜像的细粒度控制以及更稳定的幂等复制语义。用MM2做迁移本质上不是一把梭式的停机搬运而是让新旧集群短暂共存、数据实时追平随后通过切换生产者和消费者完成流量全量割接。这套方案在线上生产环境是可复现、可验证的下面我把整个实施链路拆开讲。1. 迁移场景评估与方案选型1.1 什么时候该用MM2什么时候不该用先想清楚一个问题你手里的迁移任务到底属于哪种类型我归纳过几种常见形态它们的解法完全不同。第一种是同构替换比如旧集群版本是2.13新集群是3.6变更不涉及大版本兼容性问题broker数量变化也不大。这种场景用MM2做增量同步配合一次短暂切换即可风险很低。第二种是跨集群搬迁比如从A机房迁到B机房两个集群之间物理隔离网络可能还要经过专线或公网中转。此时MM2的价值就体现出来了它有独立的连接器connector运行在Connect集群上可以单独部署在靠近源集群或目标集群的机器上灵活应对网络拓扑。第三种是版本跨度非常大的升级比如从0.10时代直接跳到3.x。这类场景我建议把迁移拆成两步先做版本升级再做数据迁移或者用MM2同步数据的同时把业务侧的序列化协议兼容问题一并解决。MM2本质上不关心消息payload格式它只负责搬运Record所以只要新旧集群对客户端协议的兼容性处理好就不存在阻碍。但也有不适合用MM2的场景。最典型的是需要精确一次且允许延迟极低的核心交易链路比如金融支付、实时风控。MM2默认是at-least-once语义即使启用幂等配置也无法像单集群内事务那样提供端到端Exactly-Once保证。这种场景我建议用双写方案由生产者主动把消息同时发给新旧集群配合消费端去重或者直接接受短暂停机做全量搬迁。另一个不适合MM2的是存量数据特别巨大且要求全量历史可查。如果Topic日志保留几天后业务就不再看历史消息那利用MM2同步最近几天的增量就够了如果要求所有历史数据都搬过去则必须先用离线工具比如kafka-dump-log配合脚本重放、或者直接用reassign工具把segment文件复制过去做基线再叠加MM2补增量。顺便提一下除了MM2现实中还有几种迁移路径。有人用Logstash或者Kafka Connect自带的JDBC源连接器做准实时抽取但那是针对业务表而非消息流。还有人用工具叫Kafka Migrator本质上是个轻量级的复制客户端也能做集群间同步但功能比MM2少很多比如没有连接器的自动故障恢复。所以我在大多数生产环境里只要集群版本在2.4以上新集群和旧集群之间的协议兼容性ok我都优先选MM2。它最大的好处是随Kafka发行版一起提供不引入额外依赖而且配置模型是声明式的容易纳入CMDB和IaC管理。1.2 MM2的核心机制与迁移思路MM2的设计思路可以这样理解它是一组特殊的Kafka Connect连接器。至少会启动两个连接器——MirrorSourceConnector负责把源集群的Topic数据复制到目标集群MirrorCheckpointConnector负责周期性同步消费者组的消费偏移。还有一个MirrorHeartbeatConnector是可选用的它会在集群间发送心跳消息用于检测链路延迟和连通性。在迁移场景中集群的角色可以互换。假设旧集群叫old新集群叫new。我们要把消息从old同步到new那么对MM2来说old是源new是目标。MM2默认会把目标集群上的Topic名称加上源集群别名作为前缀例如源集群是old同步过去的Topic名字会变成old.topic1。这一点在生产迁移时很重要因为我们通常希望迁移过程中新集群内的Topic名字保持和旧集群一致这样消费者代码不需要改名。解决办法有两种一是把新集群的集群别名配置得和旧集群一样这不符合逻辑因为别名唯一二是在MM2的配置中显式关闭topic前缀或者在做流量切换前用RenameTopic的方式把前缀去掉。这里我更推荐的是别把全量复制当成唯一路径而是用MM2先把数据复制到一个临时前缀等校验完成后在新集群侧创建最终的Topic再把临时Topic的数据可以用kafka的kafka-reassign-partitions或者干脆将MM2的复制目标改名。后面实操章节我会细说一种更优雅的做法用MM2的topic.filter配合config.properties里的replication.policy自定义类直接做到新集群内Topic名不带前缀。有一点必须提前强调MM2的同步偏移Checkpoint机制只同步source集群的consumer group偏移不会同步目标集群里同名消费者组的偏移。简单说如果新旧集群里消费者组都叫order-consumerMM2能把old里order-consumer的进度写到new里以old.order-consumer的形式存在但最终你切换消费者到新集群后到底从哪个位置继续消费取决于你在新集群里显式指定的earliest、latest还是读取old.order-consumer这个偏移。这个细节如果处理不好切换后消费者可能从头消费一切或者跳过大段消息。所以我想先给一个宏观的迁移套路后面各章节都围绕它展开在新集群上构建Connect集群并启动MM2让新旧集群的数据开始实时同步全量复制持续增量。在目标集群创建不带前缀的最终Topic并设置与源集群一致的分区数、副本数和消息格式。等同步进度赶上源集群的写入位点以后把生产者的producer配置切换为向新集群发送。处理消费者组的偏移确保切换后消费者从源集群的相同消费位置继续。业务侧逐步把消费客户端切到新集群验证消息正确性、顺序性和延迟后下线MM2和旧集群。2. 迁移前的集群准备与规划2.1 版本兼容性与网络连通性检查很多人在做MM2迁移时容易忽略版本兼容问题。MM2的源码是随Kafka发行版一起发布的不同大版本之间协议可能不兼容。Kafka的通信协议兼容性比较讲究客户端协议一般向前向后兼容几个小版本但broker间协议inter-broker protocol和connect内部的转换格式则需对齐。我的经验是如果新旧集群都是2.4及以上的版本MM2通常可以跨两个大版本同步例如从2.8同步到3.5问题不大。但如果从1.x直接同步到3.x我建议先用Kafka官方的kafka-reassign-partitions把segments直接搬迁或者用独立部署的Connect集群指定不同版本的Kafka客户端库来规避兼容问题。实操上在启动MM2前我至少会做三件检查检查新旧集群的inter.broker.protocol.version确认源和目标都能处理Kafka Connect发出的RecordBatch。从部署MM2的服务器上分别用kafka-broker-api-versions.sh --bootstrap-server 旧集群地址和指向新集群的命令验证连通性顺便确认API版本列表里包含OffsetForLeaderEpoch、Fetch等关键API。检查两个集群的message.format.version和log.message.timestamp.type避免复制消息时因为版本差异导致时间戳或格式被改写。网络方面MM2所在机器必须同时能访问新旧集群所有broker的9092内网端口以及能够解析所有broker的hostname。很多人在这里吃过亏新集群在云上旧集群在自建机房MM2部署在中间跳板机跳板机只能访问某个broker的DNAT地址却访问不了实际内网IP。结果MM2启动后连接器一直报UnresolvedAddressException。所以前置检查时必须用一个Python脚本或者nc逐一对broker列表做TCP连通性测试同时确保advertised.host.name配置的对端可路由地址。2.2 Topic映射、分区数与数据量评估在动手配置MM2之前需要先盘点源集群有哪些Topic需要迁移。重点看这几项Topic的分区数、副本数、以及现有数据总量用kafka-log-dirs.sh或du -sh统计segment目录大小。Topic的消息保留策略。如果retention.ms只有1小时那么MM2即使全量同步也只会同步最近一个小时的存活数据这在迁移中是可以接受的但如果是关键业务Topic建议临时调大保留时间比如改成7天给迁移留足缓冲。Topic的key和value的序列化类型。MM2本身不解析消息内容但你如果计划在目标集群做数据校验就需要知道是Avro、Protobuf还是纯String。我一般建议在切换前用控制台消费者把旧集群的若干消息dump出来做个样本存档万一校验出问题还能对比。这里给出一份我常用的Topic盘点表格模板项目旧集群值新集群期望值备注topic名称order-eventsorder-events前后一致分区数1212不能少于旧集群否则并行消费能力下降副本数33看broker数retention.ms172800000172800000迁移期建议扩大cleanup.policydeletedelete若compact需额外验证min.insync.replicas22保证持久性单分区日增数据量约20GB约20GB用于预估同步吞吐基于这些盘点可以估算迁移的同步窗口。比如待迁移总量是2TB网络带宽是500Mbps最理想速度约62.5MB/s算上协议开销和MM2单任务并发限制实际可能只有30MB/s。那么理论全量传输需要约19小时。如果业务允许在这个窗口内保持旧集群持续生产MM2会一边复制存量一边复制增量只要增量速度小于同步能力就能追平。追平的标志是MM2的MirrorSourceConnector的offset和源集群日志的end offset之间差得很小且稳定。有一个常见误区是认为MM2同步topic是设置好连接器后自动复制所有topic的配置。实际上MM2默认会同步所有匹配到的topic也会复制topic的配置通过内部配置Topic追踪但有些配置是不会原样复制的比如min.insync.replicas如果大于目标集群的副本数就会复制失败。我建议在启动MM2前先手动在新集群创建好所有目标Topic并把关键配置设定好再让MM2只负责数据同步不要让它自动创建。这样可控性最好。2.3 新旧集群的元数据兼容准备就算版本兼容新旧集群的某些元数据也可能存在不一致。这里指的不仅是topic配置还包括broker的rack信息、机架感知以及log.dirs路径。Kafka在收消息时关注的是分区副本的分配路径MM2复制消息时并不会重建分区的分布而是依赖目标集群的既定布局。所以在目标集群建topic时应该用kafka-topics.sh --create或AdminClient通过指定--replica-assignment来精确控制分区副本的物理分布尤其在多机架环境下避免所有副本全挤在同一机架。新旧集群如果存在password、ssl.keystore等安全配置也需要提前对齐。通常老集群可能用的是SASL_PLAINTEXT新集群切到了SASL_SSL此时MM2连接源和目标需要分别配置不同的认证机制。一个常见的坑是MM2的配置文件里可以把source.cluster.bootstrap.servers和target.cluster.bootstrap.servers设置为带SASL_SSL://前缀的地址但如果你在同一个字符串里写多个boootstrap server认证机制必须一致。否则其中一个broker返回SaslAuthenticationException整个连接器区块就会失败。另外如果源集群开启了ACL记得给MM2使用的连接器账号授予所需的Describe权限比如DescribeConfig和Read。MM2要为读取topic metadata和消费消息需要TOPIC_READ、GROUP_READ和CLUSTER_ACTION权限具体看Kafka版本。在实际配置中我倾向于在源集群为MM2创建专用账号在目标集群也创建一个只能写目标集群topic的账号避免使用超级管理员账号防止误操作。3. 核心实操MM2配置与Connect集群搭建3.1 Kafka Connect集群的独立部署MM2的运行环境是Kafka Connect可以以standalone模式或分布式模式运行。生产环境迁移我强烈建议用分布式模式因为它具备连接器自动均衡和故障恢复能力。不过这里有个细节Connect集群本身也需要一个Kafka集群来存储配置、offset和状态。在迁移早期有人图省事直接把Connect集群的internal topics放到旧集群里但这样做有两个风险一是旧集群马上要下线Connect依赖的offset topic在旧集群中一旦搬迁就会丢失连接器进度二是源目标集群的双写访问会增加旧集群的负载。正确的做法是在新集群内为Connect内部Topic单独准备一套独立的Topic。你可以直接把Connect集群的config.storage.topic、offset.storage.topic、status.storage.topic都指向新集群的同名Topic并确保副本数和分区数足够。Connect集群本身可以部署在独立的ECS或容器中不需要和Kafka broker同机。部署Connect集群时下载与新版Kafka相同版本的二进制包或镜像。如果是用Confluent Platform直接用confluent connect distributed命令即可如果是Apache Kafka发行版bin/connect-distributed.sh是默认入口。用systemd或者supervisor守护不要用nohup裸跑。我给一个简化的systemd unit示例[Unit] DescriptionApache Kafka Connect Afternetwork.target [Service] Typesimple Userkafka ExecStart/opt/kafka/bin/connect-distributed.sh /etc/kafka/connect-distributed.properties Restarton-failure RestartSec10 LimitNOFILE65536 [Install] WantedBymulti-user.targetconnect-distributed.properties里关键配置如下bootstrap.serversnew-cluster-broker-1:9092,new-cluster-broker-2:9092 group.idmm2-connect-group key.converterorg.apache.kafka.connect.converters.ByteArrayConverter value.converterorg.apache.kafka.connect.converters.ByteArrayConverter key.converter.schemas.enablefalse value.converter.schemas.enablefalse offset.storage.topicconnect-offsets config.storage.topicconnect-configs status.storage.topicconnect-status offset.storage.replication.factor3 config.storage.replication.factor3 status.storage.replication.factor3 rest.port8083 plugin.path/opt/kafka/plugins注意key.converter和value.converter在纯消息复制场景下必须设为ByteArrayConverter。很多教程直接默认用JSONConverter那样会导致消息payload被转换成JSON字符串数据内容被完整破坏。MM2这种场景本质上是字节复制不应该用带schema的converter。这也是MM2配置里比较容易踩错的点。3.2 MM2主配置文件的逐项详解在connect-distributed.properties所在目录外我们创建一个mm2.properties文件作为迁移专用配置。Kafka的MM2支持用-语法将源集群和目标集群串起来一个典型的双集群配置长这样# mm2.properties clusters old, new old.bootstrap.servers old-broker-1:9092,old-broker-2:9092 new.bootstrap.servers new-broker-1:9092,new-broker-2:9092 old-new.enabled true old-new.topics order-events, payment-events, user-activities old-new.groups order-consumer, payment-consumer old-new.emit.heartbeats.enabled true old-new.sync.topic.configs.enabled true old-new.sync.topic.acls.enabled false old-new.replication.factor 3 old-new.checkpoint.interval.ms 5000 old-new.refresh.topics.interval.ms 60000 old-new.producer.override.acks all old-new.producer.override.linger.ms 100逐项解释下关键参数clusters所有集群别名的列表用逗号分隔。后续所有和集群相关的属性都以别名为前缀。old.bootstrap.servers与new.bootstrap.servers两个集群的broker地址。如果开启SSL或SASL认证需要写完整的安全协议前缀。old-new.enabled开启从old集群到new集群的方向复制。如果还要支持回切可以考虑配new-old.enabledtrue但迁移阶段我建议先单向减少故障面。old-new.topics哪些topic需要同步。支持正则比如.*就同步所有。但在迁移中我倾向列白名单避免把内部topic和不需要的topic全部同步过去。old-new.groups需要同步偏移的消费组列表。同样可以用正则。这里建议精确列出业务消费组名单。old-new.emit.heartbeats.enabled让连接器周期性地在新集群写心跳topic用于延迟监控。可开可不开开着用于后续观察同步延迟。old-new.sync.topic.configs.enabled是否将源topic的配置同步到目标topic。如果你手动建好了目标topic并设置了正确的配置这里建议设为false避免源topic的某些配置比如本不应该同步的segment.bytes被强行复制。old-new.sync.topic.acls.enabled同步ACL。迁移阶段一般关闭。old-new.replication.factorMM2在新集群创建topic时使用的副本因子。注意这个值只对MM2自动创建topic生效手动建的Topic不受影响。通常和目标集群默认一致。old-new.checkpoint.interval.msCheckpoint连接器写偏移的时间间隔。我一般设5秒太短会增加目标集群写入压力太长则切换时可能丢失几秒的消费进度。old-new.refresh.topics.interval.ms重新扫描源集群topic列表的时间间隔。如果计划动态添加新topic就设短一点迁移阶段60秒合适。old-new.producer.override.*允许覆写目标集群生产者的配置比如acksall确保写入不丢。还有一个非常重要但容易被忽略的参数old-new.offset.syncs.topic.replication.factor 3 old-new.checkpoints.topic.replication.factor 3 old-new.heartbeats.topic.replication.factor 3这三个参数控制MM2内部创建的offset-syncs、checkpoints、heartbeats Topic的副本数。如果不显式设置默认用的是broker的默认副本因子通常也是3但有些环境默认是1这个必须盯一下。3.3 自定义复制策略去掉Topic前缀前面提到MM2默认会在目标Topic名前加源集群别名这个前缀在迁移场景里比较碍事。Kafka提供了ReplicationPolicy接口默认实现是DefaultReplicationPolicy它会用sourceAlias.originalTopic格式命名目标topic。要改变这个行为可以实现一个自定义策略类重写formatRemoteTopic方法。我用过两种办法。第一种是打开old-new.enabled并同时设置replication.policy.classorg.apache.kafka.connect.mirror.CustomReplicationPolicy然后在CustomReplicationPolicy里返回topic本身而不加前缀。这种方式的缺点是如果同时存在多个源集群指向同一个目标集群会把不同源的同名topic覆盖。但在单一迁移场景下是安全的。第二种办法我更喜欢不需要写代码先按默认策略同步让MM2把数据复制到old.order-events然后在目标集群使用kafka-topics.sh --alter重命名或者干脆在切换前用MirrorMaker的rename逻辑做一次临时topics的最终落地。实际更省事的方案是手动在新集群创建order-events通过kafka-reassign-partitions配合脚本把old.order-events的数据转移到最终topic。不过这种操作对运维要求高还容易出错。更推荐的做法是使用MM2的分区筛选和重命名设置直接把目标topic名映射成不带前缀。Kafka 2.8之后可以在MM2的配置里指定old-new.replication.policyorg.apache.kafka.connect.mirror.IdentityReplicationPolicyIdentityReplicationPolicy从Kafka 3.0开始内置支持了它唯一的作用就是使远程topic名称与源topic完全一致。但要注意这个策略下如果旧集群的主题名和目标集群中已存在的主题名重名那等同于直接往同名Topic里写这在迁移场景恰恰是我们想要的。使用IdentityReplicationPolicy时最好把old-new.sync.topic.configs.enabled设为true确保目标Topic配置自动跟随源。不过我依然建议手动预创建目标Topic避免某些配置不一致。下面是配置示例old-new.replication.policyorg.apache.kafka.connect.mirror.IdentityReplicationPolicy old-new.sync.topic.configs.enabledtrue old-new.sync.topic.acls.enabledfalse注意IdentityReplicationPolicy是在MM2的独立配置中每个方向设置确切写法是old-new.replication.policy...。如果你的Kafka版本是2.8之前的则只能自定义实现。3.4 启动MM2连接器的确认流程写完配置文件后启动Connect然后通过REST API加载MM2连接器。我这里给出一套完整的可执行步骤。首先确认Connect进程状态正常curl -s http://localhost:8083/ | jq应该返回类似{version:3.6.0,commit:...,kafka_cluster_id:...}的JSON。如果返回错误检查Connect日志。然后把mm2.properties提交为连接器任务。MM2可以作为一个自定义连接器加载也可以直接把配置文件放在Connect的classpath中用bin/kafka-mirror-maker2.sh启动。我推荐用Connect的REST API方式这样便于管理多个连接器curl -X PUT http://localhost:8083/connectors/mm2-migration/config \ -H Content-Type: application/json \ -d { connector.class: org.apache.kafka.connect.mirror.MirrorMaker2, tasks.max: 4 }但这个方式真正加载MM2的方式需要在请求体中携带所有MM2属性。我更常用的做法直接把mm2.properties转为REST请求调用PUT /connectors/mm2/config内容是一个大JSON。也可以用下面的方式cat mm2.properties | sed s/^/ /; s//: /; s/$/,/ mm2-payload.json # 手工调整成合法JSON后执行 curl -X PUT http://localhost:8083/connectors/mm2/config \ -H Content-Type: application/json \ --data mm2-payload.json提交成功后通过以下命令确认连接器所有task的状态curl -s http://localhost:8083/connectors/mm2/status | jq重点关注tasks里每个task的state字段是否为RUNNING。如果出现FAILED查看该task所在worker的日志。通常问题集中在认证失败、网络不通、Topic创建权限不足这几类。此时可以进入新集群查看topic列表应该能看到MM2创建的心跳Topic和offset-syncs Topic等比如mm2-offset-syncs.old.internal、heartbeats等。同时业务topic如果开了自动创建或手动预创建开始出现数据增长。4. 迁移实施从同步到切换的完整链路4.1 第一轮全量同步和延迟观测MM2启动后连接器的Source任务会从源集群各分区的当前起始offset开始消费。这里有个关键点MM2的Source连接器默认是自动从每个partition的earliest或latest开始消费取决于目标topic的auto.offset.reset设置以及MirrorSourceConnector内部维护的offset。如果你手动建了目标topic且没设置cleanup.policy默认情况下新消费者组从latest开始MM2可能只复制启动后的增量而不会搬存量数据。这是个容易踩坑的地方。要避免这个问题可以在MM2配置中设置old-new.poll.timeout.ms 1000 old-new.consumer.auto.offset.reset earliest其中consumer.auto.offset.resetearliest能确保Source连接器在目标topic里没有offset记录时从源分区最早位置开始读。这样就能完成全量增量同步。同步开始后务必把头30分钟作为观察期。我在这个阶段会做几件事看Connect的Worker日志中关于MirrorSourceConnector的吞吐量指标。默认Connect会暴露JMX指标可以通过jmxterm或者Grafana接入Prometheus监控。使用kafka-consumer-groups.sh在新集群查看MM2 Source任务的消费组进度。Source任务的group名称通常是mm2-mirror-old-to-new之类的内部group可以看到lag多少。如果lag持续增长说明复制速度跟不上生产速度这时候要么调大tasks.max要么调整网络和producer参数。用kafka-run-class.sh kafka.tools.GetOffsetShell对比新旧集群每个partition的最新offset。理想状态是两边offset差值越来越小直到几乎为0或只差最新几条。这里给出一个我常用的脚本片段用来统计延迟差距# 获取旧集群topic每个分区的最新offset kafka-get-offsets.sh --bootstrap-server old-broker:9092 --topic order-events old_offsets.txt # 获取新集群对应topic每个分区的最新offset kafka-get-offsets.sh --bootstrap-server new-broker:9092 --topic order-events new_offsets.txt # 用awk比较 awk NRFNR {old[$1]$2; next} {print $1, old[$1]-$2} old_offsets.txt new_offsets.txt | head -20如果差值在一段时间内稳定在一个极小值例如几十条说明已经追平。4.2 消费偏移的同步与验证MM2里负责偏移同步的是MirrorCheckpointConnector。它会周期性把源集群消费组的offset记录到目标集群的checkpoints topic中。当你在新集群里显式指定消费组从某个checkpoint读取时使用AdminClient的alterConsumerGroupOffsets或者让消费者组设置group.id不变并利用MM2的MirrorCheckpointConnector内部机制自动初始化偏移。不过实际切换时我更建议手动控制偏移。步骤是这样的在源集群记录当前所有消费组的最新提交偏移。使用kafka-consumer-groups.sh --bootstrap-server old-broker:9092 --describe --group order-consumer输出结果中包含每个partition的CURRENT-OFFSET和LOG-END-OFFSET。在新集群创建同名的消费者组可以先启动一个假消费者设置auto.offset.resetearliest然后立刻关闭这样组就存在了。使用kafka-consumer-groups的--reset-offsets命令把偏移设置为与源集群记录一致。一个常见做法是导出源集群偏移为CSV再导入新集群。但更简单的方式是用Kafka自带的kafka-consumer-groups.sh --to-earliest或者--shift-by不过这不精确。所以我写了如下Java或Python脚本用AdminClient直接修改。from kafka import KafkaAdminClient, KafkaConsumer from kafka.structs import TopicPartition # 读取旧集群offset映射 offset_map { TopicPartition(order-events, 0): 100, TopicPartition(order-events, 1): 200, } admin KafkaAdminClient(bootstrap_serversnew-broker:9092) admin.alter_consumer_group_offsets( group_idorder-consumer, offsetsoffset_map )注意admin.alter_consumer_group_offsets在较新的kafka-python版本里有如果版本低可以用kafka-python的commit_offsets配合一个临时consumer。实操上可以忍受一点点重复消费所以我一般会把偏移故意设置得比源集群小几百条这样可以确保没有消息丢失只是少量重复下游幂等处理即可。还要强调一点MM2的Checkpoint即便记录了也并不会直接修改新集群里真实消费组的offset。Checkpoint只对应一个名为old.order-consumer之类的内部group。所以不能简单认为“配置了groups同步后切过去就能从上次位置继续”。我在这点上被坑过一次后来养成了手动设置偏移的习惯。4.3 生产者和消费者切换的先后顺序与技巧标准切换流程是“先切生产者再切消费者”。原因在于如果先切消费者消费者读的是新集群但生产者还在写旧集群数据源就断了一半切完生产者后新集群的数据就持续完整了这时再切消费者就能从新集群中既读到旧数据又读到新数据。但这里有个细节切生产者前要保证新集群的Topic配置和旧集群完全一致尤其是分区数。如果分区数不一致比如旧集群12分区新集群我建成了12分区一样这还好办如果建错了分区数消费者按分区处理就可能乱掉。所以我每次在创建新集群topic前都用脚本比较两个集群的topic metadata。切换生产者的落地办法取决于客户端语言。如果是Java直接把Producer的bootstrap.servers配置从旧集群地址改成新集群地址灰度发布重启服务即可。如果用Kafka Streams或者消息中间件转发就在Streams应用配置里修改源和目的。切完生产者之后用下面这个命令监控新集群每个topic的写入速率kafka-run-class.sh kafka.tools.JmxTool --object-name kafka.server:typeBrokerTopicMetrics,nameBytesInPerSec看到新集群的BytesInPerSec在增长说明新消息已经落到新集群。消费者切换要更谨慎。如果是多个业务模块我建议按消费者组逐个切换不要一下子全切。每个消费者组切换前确认以下事项下游消息消费者是否做了幂等切换期间可能因为offset设置导致少量重复。是否需要保留旧集群的消费者组一段时间如果有人还在用旧集群旧集群不能立刻停。是否有跨consumer group的流式计算如Kafka Streams计算如果有切换必须整体协调不能只换部分。在消费者切换时最好设定消费者group的client.id和原来保持一致这样在监控图表上容易对应。如果集群之间的身份认证体系不同记得同步调整ACL和SASL配置。4.4 数据校验与流量切换验证完成切换后必须做一轮完整数据校验不能只看offset追平就觉得万事大吉。我通常分三层验证第一层数量校验。取一个时间窗口比如切换前5分钟分别统计新旧集群中该窗口内消息总量。用时间戳过滤统计因为消息中的timestamp字段在复制后会保持一致默认MM2不会重写消息时间戳除非你设置了preserve.message.timestamps为false所以统计起来相对容易。我常用Spark或者Kafka Streams做一个离线批处理也可以直接用控制台消费者消费指定时间范围的消息数。第二层内容校验。在新旧集群中分别按key抽样比对value的字节是否一致。最简单的做法是使用Python写一个校验脚本从两边topic里各自消费N条消息比如每分区1000条比较key和value的hash。需要注意消息在复制过程中不会改变序列化字节所以hash应该完全一致。如果出现不一致多半是Converters配置错误或新集群Topic出现压缩方式不同导致读取时的字节被解压重写。下面是我写过的校验脚本片段from kafka import KafkaConsumer def sample(topic, brokers, partition, count100): consumer KafkaConsumer( topic, bootstrap_serversbrokers, auto_offset_resetearliest, enable_auto_commitFalse, consumer_timeout_ms5000 ) msgs [] for msg in consumer: if msg.partition partition: msgs.append((msg.key, msg.value)) if len(msgs) count: break consumer.close() return msgs old_msgs sample(order-events, old-broker:9092, 0) new_msgs sample(order-events, new-broker:9092, 0) assert [hash((k,v)) for k,v in old_msgs] [hash((k,v)) for k,v in new_msgs]这个脚本比较粗暴没有精确对齐同一offset上的消息因为两边消费起始可能不同但用来验证内容格式是否可读、是否有大面积乱码是够用的。如果想要更严密的校验可以基于offset逐条对齐。因为源和目标topic分区的消息在MM2复制下保持原顺序和相同分区号你可以消费同一个partition的相同offset区间比对每条消息的key和value字节。比如源partition 0的offset 100到200对应目标partiton 0的offset 100到200因为MM2默认按照相同partition号复制所以可以直接按offset对齐。第三层消费链路验证。在切完消费者后观察业务日志中的消费延迟指标和监控面板确认消费者组在新集群的lag曲线从高位逐步下降说明消费链路是通的。我还会用kafka-consumer-groups.sh --describe查看所有切过来的组是不是都在正常提交位移。如果三层验证都通过我才会进入回退窗口期。保留MM2继续运行一两天但把生产者和消费者都切到新集群后其实MM2已经不再是关键路径。此时可以决定是停止MM2并清理旧集群还是再保留一段时间的双集群状态用于快速回滚。5. 迁移中的高频故障与调优实录5.1 常见错误速查表我把这些年迁移中实际遇到的典型错误整理成一个速查表看到异常日志能快速定位方向。错误现象可能原因快速排查与处理连接器状态FAILED日志出现UnresolvedAddressExceptionMM2所在机器无法解析broker hostname检查/etc/hosts或DNS从MM2机器手动nc -vz broker_ip 9092目标topic没有数据Source连接器从latest开始消费导致启动前的存量数据未被复制设置consumer.auto.offset.resetearliest并重启Source任务新集群topic名称多出前缀如old.order-events默认ReplicationPolicy导致改用IdentityReplicationPolicy或自定义策略同步时数据内容变成JSON乱码Connect的key/value converter误用JSONConverter改为ByteArrayConverter消息重复且数量远超源集群消费者组offset设置错误重复消费使用精确offset或允许少量重复并在消费端幂等连接器日志出现Topic ... not present in metadata after 60000 ms目标topic未创建且MM2创建topic权限不足手动创建topic或授权CREATE权限Checkpoint没有生成groups配置未匹配或Checkpoint task未启动检查groups正则和checkpoint.interval.ms同步过程中lag猛增单task复制能力不足、网络带宽瓶颈增加tasks.max或调节producer batch大小和压缩源和目标topic的分区数不一致导致副本错乱手动建topic时分区数配错删除重建或使用reassignment工具调整认证失败SASLAuthenticationExceptionMM2配置了错误的SASL机制或密码核对安全协议前缀与jaas配置这张表不是忽悠人的每一条我都真实遇到过。尤其是ByteArrayConverter那条有人用MM2同步后发现消息全变成了类似JSONObject结构排查半天才发现是配置了默认的JsonConverter。5.2 同步性能调优的实践参数MM2的吞吐量受多个因素影响核心瓶颈通常是单task的消费和produce速度。在分布式Connect模式下每个MirrorSourceConnector会被拆成多个task每个task负责一部分partition。tasks.max要合理设置不宜过大也不宜过小。经验值每个task建议最多负责100~200个分区。如果你有1000个分区tasks.max5~10比较合理。如果task太少单partition的最大吞吐受限太多又会导致上下文切换和内部offset管理开销增大。补几个我常用的调优参数old-new.producer.override.compression.type lz4 old-new.producer.override.batch.size 1048576 old-new.producer.override.linger.ms 20 old-new.consumer.max.poll.records 5000 old-new.consumer.fetch.min.bytes 1048576 old-new.consumer.fetch.max.bytes 52428800 old-new.consumer.max.partition.fetch.bytes 10485760注意compression.typelz4可能在某些环境存在兼容问题。如果你的目标集群消费者客户端版本较低可能不支持lz4这时改用zstd或者留空用broker默认值。消息在MM2内部默认不会重新压缩除非你设置了producer override压缩所以如果你源集群已经使用了压缩目标集群读取一样能正常解压。另一个影响性能的隐藏因素是MM2内部使用的offset.syncsTopic。MM2会周期性地记录每条记录的源offset与目标offset映射这个topic如果分区数太少在高吞吐下会成为瓶颈。建议手动创建mm2-offset-syncs.alias.internal设置至少10个分区。5.3 关于消息顺序和幂等的补充说明很多人在迁移时担心消息乱序。MM2的复制模型是按partition独立复制的MirrorSourceConnector会保证同一个源分区内的消息按原始offset顺序写入同一个目标分区。这一点做得比老版MirrorMaker好。因此只要你的业务逻辑以partition为单位保持顺序性迁移后顺序性不会丢失。但这里有两个注意点第一MM2不会保证跨partition间的全局顺序这和Kafka本身的模型一致第二如果在切换过程中让生产者在某个时间点同时写入新旧集群例如切流灰度那么同一partition的数据在不同的时间线可能分别存在于新旧集群消费者切换后不会看到同一个partition上的全局单调递增顺序而会看到一段旧数据后紧接新数据。好消息是由于MM2复制是实时的新集群的在切换时刻已经包含了旧集群存量数据的镜像所以消费者从新集群读到的是“旧数据新数据”的拼接顺序且旧数据在前的语义依然成立。幂等性方面MM2复制消息时不会改变消息的key所以目标集群的同一条消息key相同。如果下游用key做去重就能正确处理重复问题。我在切换时特意保留了旧消息的原始key同时在消费者侧增加了基于(key, partition, offset)的去重逻辑用于对冲因为offset调整可能带来的重复消费。5.4 回滚策略与常态撤离无论切换多顺利都建议准备回滚方案。我的做法是在切换完消费者后保持MM2继续运行同时让新旧两个集群都持续接收一段时间若生产者也切到新集群旧集群就没有新数据了实际上并不需要双写。最稳的回滚路径是“消费者切回旧集群生产者也切回旧集群”前提是旧集群的数据仍然完整保留着且新集群没有反向写脏数据。所以在迁移期间旧集群的retention一定要调大我一般设为至少48小时甚至7天。撤离旧集群的时机取决于所有消费者组都已不再从旧集群读取且旧集群上没有未处理的历史任务。撤离前至少观察一个完整的业务高峰周期例如24小时。如果业务侧监控显示新集群一切正常那么可以关掉MM2连接器删除Connect上的MM2任务最后再逐台下线旧集群broker。下线前务必对旧集群的server.properties和日志做个归档万一将来需要做历史追查还有据可查。6. 我总结的迁移心法说了这么多最后分享几条在实际操作中觉得最有分量的体会。第一MM2不是银弹。它对消息字节的搬运是好用的但迁移这件事的大部分工作都在外围topic配置治理、消费者组偏移管理、业务链路联排、灰度步骤设计。如果只看重配置MM2而忽略整体切换编排最后往往会在消费者切换那一刻出问题。第二把迁移当发布工程来做而不是数据工程来对待。每次迁移前我会写一份完整的checklist包含至少30个检查项从MM2机器时钟同步到新集群的日志目录挂载再到消费者组每个partition的offset对照。人脑不可靠好记性不如烂笔头。第三尽量缩短和业务侧的沟通成本。迁移期间下游消费端的负责人尤其需要知道“可能要接受少量重复消费”和“切换后消费延迟会短时波动”。最好提前给所有相关方发一份切换预案同步时间窗口和预期影响不要等出问题了再解释。最后不管你的迁移计划做了多少轮都一定要在低峰期做一次完整演练。我在生产环境切换前会在测试环境用同样的MM2配置把整个流程跑两遍一遍纯同步一遍带生产者切换和消费者偏移导入。演练完了真正切的时候手才不会抖。希望这篇基于MirrorMaker2的迁移方案能帮你在真实故障到来前把准备工作做得足够充分。
返回列表