> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Kafka Linking 最佳实践

> 规划并执行 Kafka Linking 迁移，覆盖稳定接入点、网络与权限、生产者切流、消费位点、Topic 提升、回滚和验收。

本文提供使用 Kafka Linking 将 Apache Kafka 工作负载迁移到 AutoMQ 的生产最佳实践。请先完成[前置条件梳理](/zh/automq-cloud/migrate-to-automq/prerequisites)，再结合本文制定每个迁移批次的操作手册。

<Warning>
  Kafka Linking 同步 Kafka topic 数据和 Kafka Consumer Group 已提交的消费位点。它不会改写 Flink checkpoint/savepoint、应用数据库或其他外部系统中保存的消费位点。切换消费者前，必须先确认应用实际使用的位点来源。
</Warning>

## 规划迁移批次

先选择一个非核心但具有代表性的业务完成演练，再逐步扩大范围。每个批次应包含具有依赖关系的生产者、topic 和消费者，避免把一条业务链拆到不同批次。

迁移前记录以下信息：

| 对象       | 需要确认的信息                                                     |
| -------- | ----------------------------------------------------------- |
| Topic    | 分区数、清理策略、保留时间、最大消息大小、峰值吞吐和上下游依赖                             |
| Producer | 客户端版本、是否启用幂等、是否配置 `transactional.id`、实例清单和负责人               |
| Consumer | `group.id`、客户端版本、位点来源、可接受的重复消费范围和负责人                        |
| 外部组件     | Flink、Kafka Connect、Kafka Streams、Schema Registry 和应用自管位点存储 |
| 安全配置     | 接入协议、认证信息、证书有效期、Kafka ACL、防火墙和安全组                           |
| 回滚方案     | 提升 Mirror Topic 前后的回滚步骤、暂停条件和决策人                            |

迁移窗口内避免执行源集群升级、broker 替换、网络改造、认证方式变更，以及 topic 删除重建等无关操作。减少同时发生的变更，便于判断异常来自业务切流还是基础设施变化。

## 配置稳定的源集群接入点

Kafka 客户端先通过 `bootstrap.servers` 获取集群元数据，再连接元数据中各 broker 的 `advertised.listeners` 地址。因此，只验证 Bootstrap 地址可连接并不充分。

<Warning>
  在源集群配置中填写产品提供的稳定 Bootstrap 域名或稳定的 Kafka 接入点。不要填写 Kubernetes Pod IP、节点 IP 与 NodePort 的组合，或者其他会随 Pod 重建、节点替换和扩缩容变化的临时地址。
</Warning>

按以下要求检查源集群接入：

* 使用稳定 DNS 名称或源 Kafka 服务提供的正式接入点。配置多个地址时，使其分布在不同故障域。
* 确认 AutoMQ 数据面可以解析并连接元数据中返回的**每一个** broker 地址和端口。
* 对 Kubernetes 源集群，使用 Kafka Operator 或平台提供的稳定 Bootstrap 服务和稳定的逐 broker 外部监听地址。
* 确认防火墙、安全组、路由和网络访问控制同时放行 Bootstrap 地址与全部 broker 地址。
* 使用 TLS 或 mTLS 时，确认证书信任链、域名与证书 SAN 匹配，并在迁移完成前保持证书有效。
* 使用专用迁移身份，并授予读取所选 topic、查询 Consumer Group 及完成元数据发现所需的权限。不要在迁移过程中直接轮换或删除该身份。

从与 AutoMQ 数据面网络路径等价的环境执行元数据检查。例如：

```bash theme={null}
kafka-broker-api-versions.sh \
  --bootstrap-server <stable-bootstrap-host-1>:<port>,<stable-bootstrap-host-2>:<port> \
  --command-config source-client.properties
```

命令应返回预期的 broker 列表。随后逐一验证返回地址的 DNS 解析、端口连通性和 TLS 握手。若任一 broker 地址不可达，先修复 `advertised.listeners` 或网络配置，再创建 Kafka Link。

## 准备源集群和目标实例

创建 Kafka Link 前，完成以下检查：

* 按[前置条件梳理](/zh/automq-cloud/migrate-to-automq/prerequisites)为目标实例预留迁移流量所需容量。
* 确保目标实例不存在同名 topic。Kafka Linking 需要创建与源 topic 对应的 Mirror Topic。
* 在同一个迁移批次中保持源 topic 的身份和分区结构稳定。不要删除后使用同名重新创建 topic。
* 确保目标 topic 的数据保留策略覆盖完整迁移和观察窗口。目标端不能在消费者切换前清理其仍需读取的数据。
* 单独验证目标端的 Kafka ACL、证书、Schema Registry、Connector 和其他外部依赖。不要假设 topic 数据同步会自动迁移所有外围配置和状态。
* 保存源集群与目标实例的 topic 配置、分区范围、Consumer Group 位点和业务流量基线，供切换后对比。

<Danger>
  从 Mirror Topic 创建开始，直到其完成提升并进入 `PROMOTED` 状态，不要对源 topic 或目标 Mirror Topic 扩分区。Kafka Linking 期间源、目标 topic 的分区数必须保持不变。
</Danger>

Kafka Link 进入同步状态后，持续观察复制延迟和错误。只有延迟稳定收敛且没有持续的认证、网络或请求错误时，才能开始业务切流。

<Danger>
  在全部 Mirror Topic 和 Consumer Group 完成提升且业务验收通过之前，必须保持源集群及其网络、DNS 接入点和迁移身份可用。不要提前停止、缩容到零、删除或释放源集群，也不要撤销 Kafka Linking 使用的访问权限。
</Danger>

## 切换普通生产者和消费者

本节适用于未使用 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 已完成提升，并验证消费速率、积压、业务结果和错误日志。观察窗口达到迁移计划要求后，再处理下一批业务。

<Warning>
  滚动切换必须最终收敛到全部实例连接目标端。不要让同一个 `group.id` 在源集群和目标实例上长期并行运行；这会阻止 Consumer Group 提升，并使消费进度和回滚位置难以判断。
</Warning>

## 切换事务型生产者

Kafka 事务不能通过普通生产者的滚动代理路径迁移。凡是配置了 `transactional.id`，或依赖 exactly-once 语义的应用，都按以下顺序操作：

1. 保持事务型生产者连接源集群，等待目标端复制延迟稳定收敛。
2. 停止源集群上的事务型生产者，并确认所有已开始的事务已经提交或中止。
3. 对每个分区比较源端和目标端的末端位点，确认目标端已经追平源端。
4. 提升事务涉及的全部 Mirror Topic，并等待它们进入 `PROMOTED` 状态。
5. 确认旧生产者实例不会再次启动，然后在目标实例上启动事务型生产者。
6. 使用 `read_committed` 消费验证提升前后的已提交消息均可读取，并检查事务型业务的端到端结果。

<Warning>
  Mirror Topic 仍处于 `LINKING` 状态时，不要把事务型生产者切换到目标实例。对跨多个 topic 的事务，应把所有相关 topic 纳入同一迁移批次，并统一停止、追平和提升。
</Warning>

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

消费者连接目标实例时使用哪个位点，取决于客户端的启动方式：

| 位点来源                       | Kafka Linking 的处理                      | 切换前检查                                            |
| -------------------------- | -------------------------------------- | ------------------------------------------------ |
| Kafka Consumer Group 已提交位点 | 源端实例全部退出后，Kafka Linking 同步并提升 Group 位点 | 确认每个分区的源 Group 已提交位点处于目标端可读范围内，并确保滚动发布最终退出全部源端实例 |
| Flink checkpoint/savepoint | Kafka Linking 不改写 Flink 状态             | 检查每个分区在 checkpoint/savepoint 中保存的位点              |
| 应用自管位点                     | Kafka Linking 不改写数据库、Redis、文件或业务快照     | 找到实际状态源，并检查应用传给 `seek` 的位点                       |
| Kafka Connect 或其他组件状态      | 状态可能保存在组件内部 topic 或外部存储                | 按组件的迁移流程单独验证，不要只检查业务 Consumer Group              |

对每个 topic-partition 验证：

```text theme={null}
目标端最早可读位点 <= 客户端启动位点 <= 目标端末端位点
```

如果任一分区不满足此条件，暂停该消费者切换。可选处理方式包括继续在源集群运行以推进 checkpoint、调整同步与保留范围，或在明确业务影响后重置客户端状态。

例如，Consumer 在 Partition X 的实际启动位点为 100：若目标端可读范围为 `[80, 150]`，则可以切换；若为 `[120, 150]`，说明位点 100 对应的历史消息已不可读；若为 `[80, 90]`，说明目标端尚未复制到位点 100，应继续等待目标端追平。

<Note>
  `auto.offset.reset` 只在客户端没有可用起始位点时生效。Flink restore 或应用显式调用 `assign` 和 `seek` 时，它不会替代 checkpoint 或外部存储中的位点。
</Note>

## 提升 Mirror Topic

提升会停止该 topic 的写入代理和源端数据复制。执行提升前，逐项确认以下门禁：

| 门禁             | 通过条件                                     |
| -------------- | ---------------------------------------- |
| 普通生产者          | 所有业务 Producer 实例已连接目标端；没有实例仍使用源集群接入点直接写入 |
| 事务型生产者         | 源端实例已停止；事务已结束；目标端已追平；目标端实例尚未启动           |
| Consumer Group | 源 Group 无活跃成员；目标消费者运行稳定                  |
| 外部位点消费者        | 每个分区的启动位点均在目标端可读范围内                      |
| Kafka Link     | 复制延迟已收敛，没有持续的网络、认证或请求错误                  |
| 业务验收           | 生产、消费、积压、错误率和关键业务结果符合预期                  |
| 回滚准备           | 已记录回滚负责人、数据对账方法和停止条件                     |

Kafka Linking 在 `LINKING` 阶段仍会连接源集群并代理写入，因此不能仅根据源集群存在连接判断业务 Producer 尚未切换。对低频 topic，也不要只根据短时间内流量为零判断生产者已经停止。应检查业务 Producer 实例清单和发布配置，并使用长于正常消息间隔的观察窗口对比源端和目标端写入请求。

按业务批次提升 topic，不要一次提升所有 topic。每批提升后先完成业务验证，再继续下一批。

<Note>
  提升 Mirror Topic 时，少量仍命中旧代理路径的写请求可能短暂收到 `OUT_OF_ORDER_SEQUENCE_NUMBER`。对于启用幂等性并保留正常重试配置的非事务型 Producer，Kafka 客户端通常会自动重置序列状态并重试，业务一般无需介入。建议监控 Producer 的最终发送失败；仅当错误最终暴露给应用，或在提升完成后仍持续出现时，再检查客户端配置、Topic 状态和残留代理流量，并按业务幂等策略处理明确失败的消息。
</Note>

## 规划回滚边界

提升前和提升后的回滚语义不同：

* **提升前：** 普通生产者可以切回源集群。消费者切回前需要处理源端 Consumer Group 位点，否则可能重复消费目标端已经处理的消息。
* **提升后：** 目标端的新消息不再复制到源集群。此时不能只修改 `bootstrap.servers` 直接回切，否则源集群会缺少提升后的消息。应先停止切流，并按照数据对账或反向迁移方案处理。

## 完成迁移

满足以下条件后，才能结束 Kafka Linking 迁移：

* Kafka Link 中全部 Mirror Topic 均已进入 `PROMOTED` 状态。
* Kafka Link 中全部 Consumer Group 均已完成提升。
* 目标 AutoMQ 实例上的生产、消费、积压、错误率和关键业务结果均通过验收，并完成约定的观察窗口。

确认上述条件全部满足后，在 AutoMQ Console 中直接删除 **Kafka Link 本身**。删除 Kafka Link 标志着本次 Kafka Linking 迁移彻底结束。

<Danger>
  不要逐个删除 Kafka Link 中已经提升的 Mirror Topic 或 Consumer Group。删除这些资源会真实删除目标 AutoMQ 实例中对应的 topic 或 Group，可能造成消息数据、消费位点丢失或业务中断。完成迁移时只删除 Kafka Link 本身。
</Danger>

删除 Kafka Link 后，建议继续保留源集群并静置一段观察时间：停止新的业务写入，但保留原有数据和必要的访问能力。确认目标端持续稳定且不再需要回滚或数据对账后，再回收源集群及其网络、存储、计算资源和迁移身份。完整操作步骤请参考[实施迁移方案](/zh/automq-cloud/migrate-to-automq/executing-migration)。

## 使用迁移检查表

每个批次至少保留以下记录：

* 稳定的源 Bootstrap 地址和全部 `advertised.listeners` 连通性验证结果。
* 迁移身份权限、证书有效期和防火墙放行记录。
* Topic、Producer、Consumer Group、外部位点和负责人清单。
* 切流前后的分区最早位点、末端位点、Group 已提交位点和复制延迟。
* 普通或事务型生产者所采用的切流顺序。
* 源 Group 进入 `Empty`、目标消费者稳定运行和业务验收的证据。
* Mirror Topic 提升审批、执行时间、观察结果和回滚边界。
* 全部 Mirror Topic 和 Consumer Group 的提升结果、Kafka Link 删除记录及源集群资源回收审批。

完成上述检查后，按照[实施迁移方案](/zh/automq-cloud/migrate-to-automq/executing-migration)在 AutoMQ Console 中创建 Kafka Link 并执行迁移。
