Skip to main content
本文提供使用 Kafka Linking 将 Apache Kafka 工作负载迁移到 AutoMQ 的生产最佳实践。请先完成前置条件梳理,再结合本文制定每个迁移批次的操作手册。
Kafka Linking 同步 Kafka topic 数据和 Kafka Consumer Group 已提交的消费位点。它不会改写 Flink checkpoint/savepoint、应用数据库或其他外部系统中保存的消费位点。切换消费者前,必须先确认应用实际使用的位点来源。

规划迁移批次

先选择一个非核心但具有代表性的业务完成演练,再逐步扩大范围。每个批次应包含具有依赖关系的生产者、topic 和消费者,避免把一条业务链拆到不同批次。 迁移前记录以下信息: 迁移窗口内避免执行源集群升级、broker 替换、网络改造、认证方式变更,以及 topic 删除重建等无关操作。减少同时发生的变更,便于判断异常来自业务切流还是基础设施变化。

配置稳定的源集群接入点

Kafka 客户端先通过 bootstrap.servers 获取集群元数据,再连接元数据中各 broker 的 advertised.listeners 地址。因此,只验证 Bootstrap 地址可连接并不充分。
在源集群配置中填写产品提供的稳定 Bootstrap 域名或稳定的 Kafka 接入点。不要填写 Kubernetes Pod IP、节点 IP 与 NodePort 的组合,或者其他会随 Pod 重建、节点替换和扩缩容变化的临时地址。
按以下要求检查源集群接入:
  • 使用稳定 DNS 名称或源 Kafka 服务提供的正式接入点。配置多个地址时,使其分布在不同故障域。
  • 确认 AutoMQ 数据面可以解析并连接元数据中返回的每一个 broker 地址和端口。
  • 对 Kubernetes 源集群,使用 Kafka Operator 或平台提供的稳定 Bootstrap 服务和稳定的逐 broker 外部监听地址。
  • 确认防火墙、安全组、路由和网络访问控制同时放行 Bootstrap 地址与全部 broker 地址。
  • 使用 TLS 或 mTLS 时,确认证书信任链、域名与证书 SAN 匹配,并在迁移完成前保持证书有效。
  • 使用专用迁移身份,并授予读取所选 topic、查询 Consumer Group 及完成元数据发现所需的权限。不要在迁移过程中直接轮换或删除该身份。
从与 AutoMQ 数据面网络路径等价的环境执行元数据检查。例如:
命令应返回预期的 broker 列表。随后逐一验证返回地址的 DNS 解析、端口连通性和 TLS 握手。若任一 broker 地址不可达,先修复 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 的分区数必须保持不变。
Kafka Link 进入同步状态后,持续观察复制延迟和错误。只有延迟稳定收敛且没有持续的认证、网络或请求错误时,才能开始业务切流。
在全部 Mirror Topic 和 Consumer Group 完成提升且业务验收通过之前,必须保持源集群及其网络、DNS 接入点和迁移身份可用。不要提前停止、缩容到零、删除或释放源集群,也不要撤销 Kafka Linking 使用的访问权限。

切换普通生产者和消费者

本节适用于未使用 Kafka 事务的生产者。使用 transactional.id 的生产者必须采用事务型生产者流程
  1. 将普通生产者按批次切换到目标 AutoMQ 实例,并重新创建或重启客户端。Mirror Topic 处于 LINKING 状态时,Kafka Linking 会将这些写入代理到源集群,再复制回目标实例。
  2. 检查每个生产者实例的部署配置和发送回调。确认没有未纳入迁移清单的实例继续直接写入源集群。
  3. 切换 Consumer Group 前,确认每个分区的实际启动位点均处于目标 topic 的可读范围内。不同客户端的位点来源请参考验证消费者实际使用的位点
  4. 修改 Consumer 的接入配置并滚动发布,使实例分批连接目标 AutoMQ。滚动期间源、目标集群上短暂存在同一个 group.id 属于正常切换过程;源端实例全部退出后,Kafka Linking 会完成 Consumer Group 提升。
  5. 确认全部 Consumer 实例均已连接目标端,Consumer Group 已完成提升,并验证消费速率、积压、业务结果和错误日志。观察窗口达到迁移计划要求后,再处理下一批业务。
滚动切换必须最终收敛到全部实例连接目标端。不要让同一个 group.id 在源集群和目标实例上长期并行运行;这会阻止 Consumer Group 提升,并使消费进度和回滚位置难以判断。

切换事务型生产者

Kafka 事务不能通过普通生产者的滚动代理路径迁移。凡是配置了 transactional.id,或依赖 exactly-once 语义的应用,都按以下顺序操作:
  1. 保持事务型生产者连接源集群,等待目标端复制延迟稳定收敛。
  2. 停止源集群上的事务型生产者,并确认所有已开始的事务已经提交或中止。
  3. 对每个分区比较源端和目标端的末端位点,确认目标端已经追平源端。
  4. 提升事务涉及的全部 Mirror Topic,并等待它们进入 PROMOTED 状态。
  5. 确认旧生产者实例不会再次启动,然后在目标实例上启动事务型生产者。
  6. 使用 read_committed 消费验证提升前后的已提交消息均可读取,并检查事务型业务的端到端结果。
Mirror Topic 仍处于 LINKING 状态时,不要把事务型生产者切换到目标实例。对跨多个 topic 的事务,应把所有相关 topic 纳入同一迁移批次,并统一停止、追平和提升。

验证消费者实际使用的位点

消费者连接目标实例时使用哪个位点,取决于客户端的启动方式: 对每个 topic-partition 验证:
如果任一分区不满足此条件,暂停该消费者切换。可选处理方式包括继续在源集群运行以推进 checkpoint、调整同步与保留范围,或在明确业务影响后重置客户端状态。 例如,Consumer 在 Partition X 的实际启动位点为 100:若目标端可读范围为 [80, 150],则可以切换;若为 [120, 150],说明位点 100 对应的历史消息已不可读;若为 [80, 90],说明目标端尚未复制到位点 100,应继续等待目标端追平。
auto.offset.reset 只在客户端没有可用起始位点时生效。Flink restore 或应用显式调用 assignseek 时,它不会替代 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 实例上的生产、消费、积压、错误率和关键业务结果均通过验收,并完成约定的观察窗口。
确认上述条件全部满足后,在 AutoMQ Console 中直接删除 Kafka Link 本身。删除 Kafka Link 标志着本次 Kafka Linking 迁移彻底结束。
不要逐个删除 Kafka Link 中已经提升的 Mirror Topic 或 Consumer Group。删除这些资源会真实删除目标 AutoMQ 实例中对应的 topic 或 Group,可能造成消息数据、消费位点丢失或业务中断。完成迁移时只删除 Kafka Link 本身。
删除 Kafka Link 后,建议继续保留源集群并静置一段观察时间:停止新的业务写入,但保留原有数据和必要的访问能力。确认目标端持续稳定且不再需要回滚或数据对账后,再回收源集群及其网络、存储、计算资源和迁移身份。完整操作步骤请参考实施迁移方案

使用迁移检查表

每个批次至少保留以下记录:
  • 稳定的源 Bootstrap 地址和全部 advertised.listeners 连通性验证结果。
  • 迁移身份权限、证书有效期和防火墙放行记录。
  • Topic、Producer、Consumer Group、外部位点和负责人清单。
  • 切流前后的分区最早位点、末端位点、Group 已提交位点和复制延迟。
  • 普通或事务型生产者所采用的切流顺序。
  • 源 Group 进入 Empty、目标消费者稳定运行和业务验收的证据。
  • Mirror Topic 提升审批、执行时间、观察结果和回滚边界。
  • 全部 Mirror Topic 和 Consumer Group 的提升结果、Kafka Link 删除记录及源集群资源回收审批。
完成上述检查后,按照实施迁移方案在 AutoMQ Console 中创建 Kafka Link 并执行迁移。