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 和 Consumer Group 候选列表

创建 Mirror Topic 或 Consumer Group 时,控制台从源集群 Kafka API 获取候选资源。输入搜索关键词后,控制台按名称包含关系匹配资源,例如输入 order 可以匹配 prod-order-v1,每次最多展示 100 条匹配结果。候选列表只展示 Kafka Linking 源端身份有权限查看的资源,并继续应用以下过滤规则: 如果迁移清单中的资源未显示,请按顺序检查:
  1. 缩短搜索关键词,确认资源名称能够被包含匹配,并检查是否超过单次 100 条的展示范围。
  2. 确认名称或 Group ID 没有命中上述内部资源规则。
  3. 对 Consumer Group,使用 Kafka 管理工具确认其 protocolType。非 Consumer 协议 Group 不能创建为 Kafka Linking Consumer Group。
  4. 确认 Kafka Linking 源端身份具有发现 Topic 和 Consumer Group 所需的 Kafka ACL。无权查看的资源不会出现在候选列表中。
  5. 单独核对目标实例是否已存在同名 Topic,以及该资源是否已创建为 Mirror 资源。这些冲突不会从源端候选列表中预先过滤,但可能导致创建失败。
Kafka Connect 等组件使用的内部 Topic、协调 Group 和外部状态不属于此候选列表支持的业务 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。对于使用 subscribegroup.id 参与 Group 管理的标准 Consumer,目标 Mirror Group 处于 LINKING 状态时,目标端实例无法完成 JoinGroup 和获得分区分配,因此不会开始消费。此时仍由尚未切换的源端实例承担消费流量。
  5. 控制滚动批次,并持续监控源端 Consumer 的处理能力和积压。随着源端实例逐批退出,源端有效消费能力会下降;必须确保剩余源端实例能够承载当时的全部消费负载。
  6. 确认源端 Consumer 实例全部退出。Kafka Linking 检测到源 Consumer Group 已无活跃成员后,会自动提升目标 Consumer Group。提升完成前会存在短暂的消费暂停窗口。
  7. 等待 Consumer Group 进入 PROMOTED 状态,并确认目标端 Consumer 已成功加入 Group、获得分区分配并恢复消费。验证消费速率、积压、业务结果和错误日志,观察窗口达到迁移计划要求后,再处理下一批业务。
滚动切换必须最终收敛到全部实例连接目标端。不要让同一个 group.id 在源集群和目标实例上长期并行运行;只要源 Group 仍有活跃成员,Consumer Group 就不会自动提升。目标端 Consumer 进程已启动或已连接 Broker,不代表已经开始消费;必须以 Group 进入 PROMOTED、目标 Consumer 获得分区分配并产生消费进度为准。

了解 Consumer Group 自动提升和消费暂停窗口

数据面的自动提升按以下流程执行:
  1. 目标端 Consumer 尝试加入仍处于 LINKING 状态的 Mirror Group 时,数据面登记该 Group 的自动提升检查。此次 JoinGroup 会收到可重试错误,目标端 Consumer 暂时无法获得分区分配。
  2. 数据面查询源 Group 的状态。源 Group 查无记录,或者没有活跃成员且状态为 EMPTYDEAD 时,才会尝试自动提升;如果源端仍有成员,则保持目标 Group 为 LINKING 并继续阻止目标 Consumer 加入。
  3. 满足提升条件后,Kafka Linking 获取源 Group 已提交位点,校验位点处于目标 topic 的可读范围内,将位点提交到目标 Group,并将 Group 更新为 PROMOTED
  4. 目标 Consumer 按客户端重试配置再次加入 Group,完成 rebalance 和分区分配后,才会在 AutoMQ 上继续消费。
数据面每 10 秒调度一次已登记的 Group。首次 JoinGroup 登记后,通常可在下一轮调度中发起源端检查;如果检查时源 Group 仍有活跃成员,同一 Group 的后续源端状态检查默认至少间隔 30 秒。因此,最后一个源端实例退出后,目标端恢复消费不是即时的,还需要等待下一次符合条件的源端检查、位点同步与校验、Group 提升,以及客户端重试和 rebalance。请按数十秒级消费暂停窗口制定切换计划;网络、认证、位点越界或提升失败会进一步延长该窗口。

切换事务型生产者

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 并执行迁移。