> ## 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.

# Diffusion Source Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Diffusion Source Connector，包括前置条件、配置、监控和故障排查。

## 概述

Diffusion Source Connector 订阅 Diffusion 中符合 selector 的 Topic 更新，将每次收到的值转换为 Kafka Connect SourceRecord，并写入 Kafka。它位于 Diffusion 与 Kafka 之间，适合把实时状态、设备数据或应用事件引入 Kafka，供流处理、分析和下游服务消费。

每条 Kafka 记录的键是完整的 Diffusion Topic 路径，值来自 Diffusion JSON 数据。目标 Kafka Topic 可固定为一个名称，也可使用 `${topic}` 根据 Diffusion Topic 路径动态生成；动态映射时路径中的 `/` 会替换为 `_`。Connector 采用单 Task 订阅和传输数据，不将 Diffusion Topic 自动拆分到多个 Task。

## 前置条件

* 要传输数据，Diffusion 中需要有与 selector 匹配的 JSON Topic；连接账号需要具备建立会话、选择并订阅这些 Topic 的权限。
* 目标 Kafka Topic 应满足名称规则并允许 Connector 写入；使用动态映射时，应提前检查路径转换后的 Topic 名称及可能的名称碰撞。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka、Diffusion 服务和需要订阅的 JSON Topic，确认网络连通和访问权限；创建与管理操作参见 AutoMQ 的[管理 Connector](../manage-connectors)。以下配置把匹配 `source/app/` 路径的更新汇聚到一个 Kafka Topic。

```properties theme={null}
connector.class=com.diffusiondata.connect.diffusion.source.DiffusionSourceConnector
diffusion.url=<diffusion-url>
diffusion.username=<diffusion-principal>
diffusion.password=<diffusion-password>
diffusion.selector=?source/app/.*
kafka.topic=diffusion-events
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
```

替换 Diffusion 连接地址和凭据，并确认 selector 只匹配计划接入的 JSON Topic。示例将所有匹配更新写入 `diffusion-events`，记录键仍保留原始 Diffusion Topic 路径。密码应通过受控的配置或密钥管理机制提供，不要写入源码、日志或共享材料。JsonConverter 是否在 JSON 中包含 Schema 信息，取决于运行环境中该 Converter 的配置。

## 配置

### Diffusion 连接与认证

#### `diffusion.url`

指定建立 Diffusion 会话使用的完整连接 URL。

* **类别**：Diffusion 连接
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：填写 Diffusion 客户端可接受且 Worker 能访问的 URL，例如 WebSocket 地址。配置解析不校验 URL 语法、地址可达性或空字符串，实际连接失败会使 Task 启动失败。

#### `diffusion.username`

指定连接 Diffusion 时使用的 principal 用户名。

* **类别**：Diffusion 认证
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：账号需要具备建立会话以及选择、订阅 `diffusion.selector` 所匹配 Topic 的权限。配置解析不校验账号是否存在或是否具备权限。

#### `diffusion.password`

指定 Diffusion principal 的密码。

* **类别**：Diffusion 认证
* **类型**：`password`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：通过受控的密钥管理机制提供真实密码，避免出现在源码、日志或共享材料中。密码是否有效由 Diffusion 在连接时检查。

### 订阅范围与 Topic 映射

#### `diffusion.selector`

指定订阅 Diffusion Topic 的 selector。

* **类别**：源数据选择
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：使用 Diffusion 支持的 selector 语法，并只匹配计划接入的 JSON Topic。selector 原样交给 Diffusion 解析；配置校验不检查语法、匹配结果或权限。订阅在启动阶段未能及时完成时，Task 会启动失败。

#### `kafka.topic`

指定 Diffusion Topic 路径到 Kafka Topic 名称的映射模式。

* **类别**：目标 Topic
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：固定名称会把所有匹配更新汇聚到同一 Kafka Topic；`${topic}` 会替换为完整 Diffusion Topic 路径，随后结果中的 `/` 全部替换为 `_`。其他 token 不会替换。最终名称必须符合 Kafka Topic 规则并具备写权限；不同路径可能映射为同名 Topic。

### 轮询与批量

#### `diffusion.poll.interval`

指定内存队列暂时没有记录时，每次再次检查前的等待时间，单位毫秒。

* **类别**：轮询
* **类型**：`int`
* **默认值**：`1000`
* **重要级别**：低
* **有效值 / 注意事项**：使用 `0` 至 `2147483647`。一次空轮询最多执行三次等待，因此返回空结果前的空闲等待可能接近该值的三倍。`0` 仅取消检查之间的等待，不改变固定检查次数；负值会在运行时失败。

#### `diffusion.poll.size`

指定一次 poll 最多从内存队列返回的记录数。

* **类别**：批量
* **类型**：`int`
* **默认值**：`128`
* **重要级别**：低
* **有效值 / 注意事项**：使用 `1` 至 `2147483647`。该值只限制单次返回数量，不限制队列总大小，也不是 Kafka Producer 的 `batch.size`。较大值可能增加单批转换和内存处理量；`0` 或负值不能形成可用的批量行为。

### 运行与序列化

#### `connector.class`

指定 Diffusion Source Connector 实现类。

* **类别**：Kafka Connect 框架配置
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：使用 `com.diffusiondata.connect.diffusion.source.DiffusionSourceConnector`。

#### `tasks.max`

指定 Connector 可使用的最大 Task 数。

* **类别**：Kafka Connect 框架配置
* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：至少为 `1`。此 Connector 始终只生成一个 Task，提高该值不会拆分 selector、增加 Source 并行度或改变单 Task 行为，应保留为 `1`。

#### `tasks.max.enforce`

指定 Kafka Connect 是否强制 Connector 生成的 Task 数不超过 `tasks.max`。

* **类别**：Kafka Connect 框架配置
* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **已弃用**：是
* **替代项**：无
* **有效值 / 注意事项**：可取 `true` 或 `false`。Kafka Connect 已弃用该配置，应保留默认值；关闭限制不会使此 Connector 生成多个 Task。

#### `key.converter`

指定 Diffusion Topic 路径记录键的序列化 Converter。

* **类别**：Kafka Connect 框架配置
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：省略时继承 Worker 的 Key Converter。记录键为字符串形式的完整 Diffusion Topic 路径；需要按纯文本序列化记录键时可使用 `org.apache.kafka.connect.storage.StringConverter`，该类不是默认值。

#### `value.converter`

指定 Diffusion JSON 转换后的记录值使用的序列化 Converter。

* **类别**：Kafka Connect 框架配置
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：省略时继承 Worker 的 Value Converter。Connector 会先把 JSON 值转换为 Kafka Connect 数据及动态推断的 Schema；需要 JSON 输出时可使用 `org.apache.kafka.connect.json.JsonConverter`，该类不是默认值。Converter 专属设置由所选 Converter 和运行环境管理。

## 最佳实践

### 按 Diffusion 路径拆分目标 Kafka Topic

**适用业务场景**：已使用固定 Kafka Topic 完成初始连通性验证，需要让不同 Diffusion Topic 的数据进入不同 Kafka Topic，以便下游按业务流分别设置保留策略、权限或消费任务，同时保留来源路径与目标名称之间的直观对应关系。

**配置示例**：以快速开始为基础，将 `kafka.topic` 改为以下模式，其余配置保持不变。

```properties theme={null}
kafka.topic=diffusion_${topic}
```

**关键说明**：例如 Diffusion 路径 `source/app/orders` 会映射为 `diffusion_source_app_orders`。上线前列出 selector 可能匹配的路径，检查转换后的名称是否合法且互不冲突，并准备相应 Kafka Topic 或自动创建策略与写权限。映射只替换 `${topic}` 和 `/`，不会清理其他非法字符；`foo/bar` 与 `foo_bar` 可能得到同一结果。若业务需要把所有更新汇聚到同一 Topic，应继续使用快速开始中的固定名称，不要同时追求动态拆分。

### 运行期间调整等待时间与单批数量

**适用业务场景**：Connector 已稳定传输数据，需要在空闲期新数据等待时间、单次处理量和 Worker 内存压力之间进行调整。

**配置示例**：在现有 Connector 配置中覆盖以下两项，其余连接、selector、Topic 映射和 Converter 配置保持不变。

```properties theme={null}
diffusion.poll.interval=500
diffusion.poll.size=256
```

**关键说明**：示例值只是调整起点，不是通用最优值。缩短等待时间会让空队列检查更频繁；增大单批数量允许一次返回更多已排队记录，但会增加单批转换和内存处理量。结合 Kafka 写入吞吐、端到端延迟、Worker 堆内存和 GC 调整，并持续观察积压趋势。该配置不会给内存队列增加容量上限，也不会向 Diffusion 提供背压。

## 监控

### 监控内容

关注 Kafka Connect 集群健康、Connector 与 Task 状态、记录吞吐、端到端延迟、Offset 提交、错误和重试，以及 Worker JVM 内存、GC 和线程信号；Task 处于 RUNNING 不代表 selector 持续收到更新，应结合 Diffusion 源数据变化和 Kafka 实际输出判断。由于 Connector 内部队列没有容量上限，还应特别关注持续流量下的堆内存增长。只有部署启用了相应错误处理时才关注 DLQ 活动。

### 导入 Grafana 大盘

下载共享的 [AutoMQ Connect Cluster Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确认 Grafana 数据源能够查询 Kafka Connect 和 Worker JVM 指标，并与大盘所需的集群、Connector、Task 等标签匹配，再通过 Grafana 的导入功能加载 JSON 并选择对应数据源。

## 限制条件

* 单个 Connector 实例始终只有一个 Task，`tasks.max` 不能拆分 selector 或提高并行度。
* SourceRecord 不包含可恢复的 source offset，Connector 不能从 Kafka Connect checkpoint 恢复 Diffusion 更新位置，也不能回放 Task 停止期间的更新。
* Connector 不提供可据此声明的 exactly-once、at-least-once 或无重复交付保障；业务应根据可重放性和下游幂等要求评估使用方式。
* Diffusion 回调与 poll 之间使用无界内存队列，没有源端背压或总积压上限；持续输入快于 Kafka 发送时，内存占用可能持续增长。
* 每条 JSON 消息独立推断值 Schema，数据形状变化可能改变或移除 Schema；JSON 对象不会恢复为预定义的 Connect Struct。
* Topic 映射只将 `/` 替换为 `_`，不会清理其他非法字符或检测名称碰撞。
* Connector 创建 Diffusion 会话时不启用客户端自动重连，也没有内建的重新连接和重新订阅循环。

## 常见问题

### Task 正常运行，但 Kafka Topic 中没有新消息怎么办？

先确认 Diffusion 中确实有新的 JSON Topic 值更新，并检查 `diffusion.selector` 是否匹配对应路径。再核对 principal 的选择与订阅权限、`kafka.topic` 的最终名称、Kafka 写权限以及 Key 和 Value Converter。selector 没有匹配项不会在配置校验阶段报错；使用固定 Topic 时不要误查动态派生名称，使用动态映射时应按路径转换结果查找。查看 Worker 日志中的订阅、转换和 Kafka 写入错误，但不要输出真实密码。

### 为什么不同 Diffusion Topic 的数据进入了同一个 Kafka Topic？

如果 `kafka.topic` 是固定名称，所有 selector 匹配的数据都会汇聚到该 Topic。需要按来源拆分时使用包含 `${topic}` 的模式，并提前检查转换结果。即使使用动态模式，`foo/bar` 与 `foo_bar` 也会因 `/` 替换为 `_` 而发生名称碰撞；调整源路径规划或添加不会冲突的固定前缀，并逐一核对最终 Topic 名称。

### 为什么重启后没有补回停机期间的更新？

Connector 不保存可用于恢复 Diffusion 更新位置的 source offset，重启时会建立新会话并重新订阅 selector，不能依靠 Kafka Connect Offset 回放停机期间的数据。维护前应评估源端是否能另行保留或重放所需事件，并让下游能够处理可能重复出现的值；不要把 Connector 重启视为补数机制。

### Diffusion 连接中断后 Task 没有自行恢复怎么办？

检查 Diffusion 服务、网络、账号权限和会话地址，修复故障后通过 Kafka Connect 管理界面确认 Task 状态并执行受控重启。该 Connector 没有内建的自动重连和重新订阅循环，因此不能只等待 Task 自行恢复。重启前同时评估中断期间数据无法通过 Connector offset 补回的影响。
