规划迁移批次
先选择一个非核心但具有代表性的业务完成演练,再逐步扩大范围。每个批次应包含具有依赖关系的生产者、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 和 Consumer Group 候选列表
创建 Mirror Topic 或 Consumer Group 时,控制台从源集群 Kafka API 获取候选资源。输入搜索关键词后,控制台按名称包含关系匹配资源,例如输入order 可以匹配 prod-order-v1,每次最多展示 100 条匹配结果。候选列表只展示 Kafka Linking 源端身份有权限查看的资源,并继续应用以下过滤规则:
如果迁移清单中的资源未显示,请按顺序检查:
- 缩短搜索关键词,确认资源名称能够被包含匹配,并检查是否超过单次 100 条的展示范围。
- 确认名称或 Group ID 没有命中上述内部资源规则。
- 对 Consumer Group,使用 Kafka 管理工具确认其
protocolType。非 Consumer 协议 Group 不能创建为 Kafka Linking Consumer Group。 - 确认 Kafka Linking 源端身份具有发现 Topic 和 Consumer Group 所需的 Kafka ACL。无权查看的资源不会出现在候选列表中。
- 单独核对目标实例是否已存在同名 Topic,以及该资源是否已创建为 Mirror 资源。这些冲突不会从源端候选列表中预先过滤,但可能导致创建失败。
从 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。对于使用
subscribe和group.id参与 Group 管理的标准 Consumer,目标 Mirror Group 处于LINKING状态时,目标端实例无法完成 JoinGroup 和获得分区分配,因此不会开始消费。此时仍由尚未切换的源端实例承担消费流量。 - 控制滚动批次,并持续监控源端 Consumer 的处理能力和积压。随着源端实例逐批退出,源端有效消费能力会下降;必须确保剩余源端实例能够承载当时的全部消费负载。
- 确认源端 Consumer 实例全部退出。Kafka Linking 检测到源 Consumer Group 已无活跃成员后,会自动提升目标 Consumer Group。提升完成前会存在短暂的消费暂停窗口。
- 等待 Consumer Group 进入
PROMOTED状态,并确认目标端 Consumer 已成功加入 Group、获得分区分配并恢复消费。验证消费速率、积压、业务结果和错误日志,观察窗口达到迁移计划要求后,再处理下一批业务。
了解 Consumer Group 自动提升和消费暂停窗口
数据面的自动提升按以下流程执行:- 目标端 Consumer 尝试加入仍处于
LINKING状态的 Mirror Group 时,数据面登记该 Group 的自动提升检查。此次 JoinGroup 会收到可重试错误,目标端 Consumer 暂时无法获得分区分配。 - 数据面查询源 Group 的状态。源 Group 查无记录,或者没有活跃成员且状态为
EMPTY或DEAD时,才会尝试自动提升;如果源端仍有成员,则保持目标 Group 为LINKING并继续阻止目标 Consumer 加入。 - 满足提升条件后,Kafka Linking 获取源 Group 已提交位点,校验位点处于目标 topic 的可读范围内,将位点提交到目标 Group,并将 Group 更新为
PROMOTED。 - 目标 Consumer 按客户端重试配置再次加入 Group,完成 rebalance 和分区分配后,才会在 AutoMQ 上继续消费。
切换事务型生产者
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 删除记录及源集群资源回收审批。