规划迁移批次
先选择一个非核心但具有代表性的业务完成演练,再逐步扩大范围。每个批次应包含具有依赖关系的生产者、topic 和消费者,避免把一条业务链拆到不同批次。 迁移前记录以下信息:
迁移窗口内避免执行源集群升级、broker 替换、网络改造、认证方式变更,以及 topic 删除重建等无关操作。减少同时发生的变更,便于判断异常来自业务切流还是基础设施变化。
配置稳定的源集群接入点
Kafka 客户端先通过bootstrap.servers 获取集群元数据,再连接元数据中各 broker 的 advertised.listeners 地址。因此,只验证 Bootstrap 地址可连接并不充分。
按以下要求检查源集群接入:
- 使用稳定 DNS 名称或源 Kafka 服务提供的正式接入点。配置多个地址时,使其分布在不同故障域。
- 确认 AutoMQ 数据面可以解析并连接元数据中返回的每一个 broker 地址和端口。
- 对 Kubernetes 源集群,使用 Kafka Operator 或平台提供的稳定 Bootstrap 服务和稳定的逐 broker 外部监听地址。
- 确认防火墙、安全组、路由和网络访问控制同时放行 Bootstrap 地址与全部 broker 地址。
- 使用 TLS 或 mTLS 时,确认证书信任链、域名与证书 SAN 匹配,并在迁移完成前保持证书有效。
- 使用专用迁移身份,并授予读取所选 topic、查询 Consumer Group 及完成元数据发现所需的权限。不要在迁移过程中直接轮换或删除该身份。
advertised.listeners 或网络配置,再创建 Kafka Link。
准备源集群和目标实例
创建 Kafka Link 前,完成以下检查:- 按前置条件梳理为目标实例预留迁移流量所需容量。
- 确保目标实例不存在同名 topic。Kafka Linking 需要创建与源 topic 对应的 Mirror Topic。
- 在同一个迁移批次中保持源 topic 的身份和分区结构稳定。不要删除后使用同名重新创建 topic。
- 确保目标 topic 的数据保留策略覆盖完整迁移和观察窗口。目标端不能在消费者切换前清理其仍需读取的数据。
- 单独验证目标端的 Kafka ACL、证书、Schema Registry、Connector 和其他外部依赖。不要假设 topic 数据同步会自动迁移所有外围配置和状态。
- 保存源集群与目标实例的 topic 配置、分区范围、Consumer Group 位点和业务流量基线,供切换后对比。
从 Mirror Topic 创建开始,直到其完成提升并进入
PROMOTED 状态,不要对源 topic 或目标 Mirror Topic 扩分区。Kafka Linking 期间源、目标 topic 的分区数必须保持不变。在全部 Mirror Topic 和 Consumer Group 完成提升且业务验收通过之前,必须保持源集群及其网络、DNS 接入点和迁移身份可用。不要提前停止、缩容到零、删除或释放源集群,也不要撤销 Kafka Linking 使用的访问权限。
切换普通生产者和消费者
本节适用于未使用 Kafka 事务的生产者。使用transactional.id 的生产者必须采用事务型生产者流程。
- 将普通生产者按批次切换到目标 AutoMQ 实例,并重新创建或重启客户端。Mirror Topic 处于
LINKING状态时,Kafka Linking 会将这些写入代理到源集群,再复制回目标实例。 - 检查每个生产者实例的部署配置和发送回调。确认没有未纳入迁移清单的实例继续直接写入源集群。
- 切换 Consumer Group 前,确认每个分区的实际启动位点均处于目标 topic 的可读范围内。不同客户端的位点来源请参考验证消费者实际使用的位点。
- 修改 Consumer 的接入配置并滚动发布,使实例分批连接目标 AutoMQ。滚动期间源、目标集群上短暂存在同一个
group.id属于正常切换过程;源端实例全部退出后,Kafka Linking 会完成 Consumer Group 提升。 - 确认全部 Consumer 实例均已连接目标端,Consumer Group 已完成提升,并验证消费速率、积压、业务结果和错误日志。观察窗口达到迁移计划要求后,再处理下一批业务。
切换事务型生产者
Kafka 事务不能通过普通生产者的滚动代理路径迁移。凡是配置了transactional.id,或依赖 exactly-once 语义的应用,都按以下顺序操作:
- 保持事务型生产者连接源集群,等待目标端复制延迟稳定收敛。
- 停止源集群上的事务型生产者,并确认所有已开始的事务已经提交或中止。
- 对每个分区比较源端和目标端的末端位点,确认目标端已经追平源端。
- 提升事务涉及的全部 Mirror Topic,并等待它们进入
PROMOTED状态。 - 确认旧生产者实例不会再次启动,然后在目标实例上启动事务型生产者。
- 使用
read_committed消费验证提升前后的已提交消息均可读取,并检查事务型业务的端到端结果。
验证消费者实际使用的位点
消费者连接目标实例时使用哪个位点,取决于客户端的启动方式:
对每个 topic-partition 验证:
[80, 150],则可以切换;若为 [120, 150],说明位点 100 对应的历史消息已不可读;若为 [80, 90],说明目标端尚未复制到位点 100,应继续等待目标端追平。
auto.offset.reset 只在客户端没有可用起始位点时生效。Flink restore 或应用显式调用 assign 和 seek 时,它不会替代 checkpoint 或外部存储中的位点。提升 Mirror Topic
提升会停止该 topic 的写入代理和源端数据复制。执行提升前,逐项确认以下门禁:
Kafka Linking 在
LINKING 阶段仍会连接源集群并代理写入,因此不能仅根据源集群存在连接判断业务 Producer 尚未切换。对低频 topic,也不要只根据短时间内流量为零判断生产者已经停止。应检查业务 Producer 实例清单和发布配置,并使用长于正常消息间隔的观察窗口对比源端和目标端写入请求。
按业务批次提升 topic,不要一次提升所有 topic。每批提升后先完成业务验证,再继续下一批。
提升 Mirror Topic 时,少量仍命中旧代理路径的写请求可能短暂收到
OUT_OF_ORDER_SEQUENCE_NUMBER。对于启用幂等性并保留正常重试配置的非事务型 Producer,Kafka 客户端通常会自动重置序列状态并重试,业务一般无需介入。建议监控 Producer 的最终发送失败;仅当错误最终暴露给应用,或在提升完成后仍持续出现时,再检查客户端配置、Topic 状态和残留代理流量,并按业务幂等策略处理明确失败的消息。规划回滚边界
提升前和提升后的回滚语义不同:- 提升前: 普通生产者可以切回源集群。消费者切回前需要处理源端 Consumer Group 位点,否则可能重复消费目标端已经处理的消息。
- 提升后: 目标端的新消息不再复制到源集群。此时不能只修改
bootstrap.servers直接回切,否则源集群会缺少提升后的消息。应先停止切流,并按照数据对账或反向迁移方案处理。
完成迁移
满足以下条件后,才能结束 Kafka Linking 迁移:- Kafka Link 中全部 Mirror Topic 均已进入
PROMOTED状态。 - Kafka Link 中全部 Consumer Group 均已完成提升。
- 目标 AutoMQ 实例上的生产、消费、积压、错误率和关键业务结果均通过验收,并完成约定的观察窗口。
不要逐个删除 Kafka Link 中已经提升的 Mirror Topic 或 Consumer Group。删除这些资源会真实删除目标 AutoMQ 实例中对应的 topic 或 Group,可能造成消息数据、消费位点丢失或业务中断。完成迁移时只删除 Kafka Link 本身。
使用迁移检查表
每个批次至少保留以下记录:- 稳定的源 Bootstrap 地址和全部
advertised.listeners连通性验证结果。 - 迁移身份权限、证书有效期和防火墙放行记录。
- Topic、Producer、Consumer Group、外部位点和负责人清单。
- 切流前后的分区最早位点、末端位点、Group 已提交位点和复制延迟。
- 普通或事务型生产者所采用的切流顺序。
- 源 Group 进入
Empty、目标消费者稳定运行和业务验收的证据。 - Mirror Topic 提升审批、执行时间、观察结果和回滚边界。
- 全部 Mirror Topic 和 Consumer Group 的提升结果、Kafka Link 删除记录及源集群资源回收审批。