> ## 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 Sink Connector

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

## 概述

Diffusion Sink Connector 将 Kafka Topic 中的记录写入 Diffusion Server。Connector 会把每条记录的值转换为 Diffusion 可处理的 JSON，并根据配置解析目标 Diffusion Topic 路径。目标 Topic 不存在时，Connector 会将其创建为 JSON Topic 后写入；目标 Topic 已存在时则直接更新。

Connector 可以把一个或多个 Kafka Topic 映射到固定的 Diffusion 路径，也可以使用 `${topic}`、`${key}`、`${key.version}` 和 `${value.version}`，根据记录元数据生成路径。典型场景是把 Kafka 中的业务事件发布到 Diffusion，供浏览器、移动端和 IoT 客户端实时消费。

## 前置条件

* Diffusion Server 为 6.9 或更高版本，并已准备接收数据的服务账号；该账号需要对目标路径执行更新，且在目标 Topic 不存在时具备创建 JSON Topic 的权限。
* 运行环境使用 Java 11 或更高版本，并能通过 `diffusion.url` 访问 Diffusion Server。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和 Diffusion Server，并确认网络连通及服务账号权限。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。

```properties theme={null}
connector.class=com.diffusiondata.connect.diffusion.sink.DiffusionSinkConnector
topics=price
diffusion.url=ws://<diffusion-host>:8080
diffusion.username=<diffusion-username>
diffusion.password=<diffusion-password>
diffusion.destination=kafka/${topic}
```

将 `<diffusion-host>`、`<diffusion-username>` 和 `<diffusion-password>` 替换为实际连接信息。此配置会把 Kafka Topic `price` 中的记录写入 Diffusion Topic `kafka/price`。提交配置时只需使用上述 Connector properties，不要包含 REST 请求外壳。

## 配置

### Diffusion 连接

#### `diffusion.url`

Diffusion Server 的连接地址。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填。使用运行 Kafka Connect Worker 可访问的完整地址；Connector 配置本身只按字符串解析，不校验协议、主机、端口或内容是否为空。
* **必填**：是

#### `diffusion.username`

用于向 Diffusion Server 认证的 principal。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填。账号必须具备访问目标 Diffusion Topic 的权限；Connector 配置本身不校验内容是否为空。
* **必填**：是

#### `diffusion.password`

用于向 Diffusion Server 认证的密码。

* **类型**：`password`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填。请使用 Connect 的敏感配置管理方式保存，不要将密码写入日志或公开配置仓库。
* **必填**：是

### 目标路径

#### `diffusion.destination`

生成 Diffusion Topic 路径的模式。Connector 会针对每条 SinkRecord 解析该模式。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填。支持 `${topic}`、`${key}`、`${key.version}` 和 `${value.version}`。`${topic}` 使用输入 Kafka Topic；其余令牌依赖记录键或 Schema 版本。记录键为空时，Connector 在替换 `${topic}` 后不会继续解析键和版本令牌；没有可用值的令牌会保留为字面量。路径分隔符不会被清理或替换。
* **必填**：是

### 输入订阅

#### `topics`

要消费的 Kafka Topic 列表。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：高
* **有效值 / 注意事项**：与 `topics.regex` 二选一。使用逗号分隔的 Topic 名称；至少配置一个非空 Topic。Topic 名称可通过 `${topic}` 参与目标路径生成。

#### `topics.regex`

按 Java 正则表达式订阅 Kafka Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：与 `topics` 二选一。必须是可编译且非空的 Java 正则表达式；如果配置了 DLQ Topic，正则不能匹配该 DLQ Topic。

### Connector 身份与任务

#### `connector.class`

要加载的 Sink Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：使用 `com.diffusiondata.connect.diffusion.sink.DiffusionSinkConnector`。
* **必填**：是

#### `tasks.max`

该 Connector 允许使用的最大 Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须大于或等于 `1`。该实现始终返回一个 Task 配置，因此提高此值不会让单个 Connector 实例创建多个 Diffusion Sink Task。

#### `tasks.max.enforce`

是否启用 Kafka Connect 对 `tasks.max` 的框架约束。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：低
* **有效值 / 注意事项**：可选。此配置在 Kafka Connect 3.9.1 中已弃用且没有替代项；它不会改变该 Connector 始终创建一个 Task 的行为。
* **已弃用**：是

### 数据转换

#### `key.converter`

指定 Kafka 记录键的 Converter。转换后的键可用于 `${key}`，其 Schema 版本可用于 `${key.version}`。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：可选。省略时继承 Worker 的键 Converter；配置时必须是可实例化的具体 Converter 类。

#### `value.converter`

指定 Kafka 记录值的 Converter。转换后的值会被序列化为 Diffusion JSON，其 Schema 版本可用于 `${value.version}`。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：可选。省略时继承 Worker 的值 Converter；配置时必须是可实例化的具体 Converter 类。Map 的键需要能够表示为字符串；不符合要求的值会导致记录处理失败。

### 错误处理

#### `errors.retry.timeout`

可重试错误的总重试时长，单位为毫秒。

* **类型**：`long`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：`0` 表示禁用重试，`-1` 表示无限重试，其他非负值表示重试时长。该配置是 Kafka Connect 的通用策略，不替代 Sink Task 等待 Diffusion 发布结果的固定等待行为。

#### `errors.tolerance`

发生错误时允许继续处理的范围。

* **类型**：`string`
* **默认值**：`none`
* **重要级别**：中
* **有效值 / 注意事项**：可选值为 `none` 或 `all`。该配置是通用错误策略；仅配置它并不能保证 Diffusion 发布失败会自动写入 DLQ。

#### `errors.deadletterqueue.topic.name`

用于保存错误记录的 Kafka DLQ Topic 名称。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：空字符串表示不发布到 DLQ。配置非空值时，该 Topic 不能同时出现在 `topics` 中，也不能被 `topics.regex` 匹配。

#### `errors.deadletterqueue.topic.replication.factor`

由 Kafka Connect 创建 DLQ Topic 时使用的副本因子。

* **类型**：`short`
* **默认值**：`3`
* **重要级别**：中
* **有效值 / 注意事项**：仅在配置的 DLQ Topic 不存在且由 Kafka Connect 创建时使用；该值必须适合 Kafka 集群的副本数量。

#### `errors.deadletterqueue.context.headers.enable`

是否在写入 DLQ 的记录中添加错误上下文 Header。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：仅在实际产生 DLQ 记录时生效；启用后错误上下文使用 `__connect.errors.` 前缀。

## 最佳实践

### 按 Topic 规则接入一组 Kafka 数据流

**适用业务场景**：需要持续接入一组名称遵循统一规则的 Kafka Topic，并把每个 Topic 的记录写入对应的 Diffusion 路径。此时使用正则订阅可以避免逐个维护 Topic 列表。

**配置示例**：

```properties theme={null}
connector.class=com.diffusiondata.connect.diffusion.sink.DiffusionSinkConnector
topics.regex=orders-.*
diffusion.url=ws://<diffusion-host>:8080
diffusion.username=<diffusion-username>
diffusion.password=<diffusion-password>
diffusion.destination=events/${topic}
```

**关键说明**：`topics.regex` 非空时不要同时配置非空的 `topics`。`${topic}` 会保留输入 Topic 名称，因此 `orders-created` 会写入 `events/orders-created`。订阅范围扩大前，应确认 Diffusion principal 对所有可能生成的目标路径都有更新或创建权限。

### 用多个 Connector 实例分隔数据流

**适用业务场景**：需要分开管理不同业务流，分别设置目标路径或故障边界，且单个 Connector 实例中的 Topic 集合已经较大。此时可为每组 Topic 创建独立的 Connector 实例。

**配置示例**：

```properties theme={null}
connector.class=com.diffusiondata.connect.diffusion.sink.DiffusionSinkConnector
topics=orders,inventory
tasks.max=1
diffusion.url=ws://<diffusion-host>:8080
diffusion.username=<diffusion-username>
diffusion.password=<diffusion-password>
diffusion.destination=business/${topic}
```

为另一组 Topic 创建第二个 Connector 实例时，应使用不同的 `topics` 或 `topics.regex` 和目标路径，而不是只提高同一实例的 `tasks.max`。该实现始终只创建一个 Task，增加 `tasks.max` 不会在实例内部产生任务级并行。多个实例之间应明确划分输入 Topic，避免重复消费同一数据流。

## 监控

### 监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号；仅在启用了相应错误处理时关注 DLQ 活动。

### 导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源，且采集标签满足大盘筛选条件；下载 [Kafka Connect Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* 单个 Diffusion Sink Connector 实例始终只创建一个 Task；提高 `tasks.max` 不会让该实例并行创建多个 Task，需要通过多个实例分隔输入 Topic 才能进行部署级扩展。
* `topics` 和 `topics.regex` 必须且只能选择一个非空订阅方式；两者同时非空或同时为空都会导致 Sink 配置校验失败。
* Diffusion 更新成功后，如果在框架提交 Kafka Offset 前发生连接丢失或 Task 重启，Offset 可能未提交；恢复后记录可能再次写入同一路径，因此 Connector 不保证应用层副作用只发生一次。
* 记录键为空时，`diffusion.destination` 只替换 `${topic}`，`${key}`、`${key.version}` 和 `${value.version}` 不会继续解析；因此依赖这些令牌的路径模式应确保输入记录具有所需的键和 Schema 版本。

## 常见问题

### 为什么增加 `tasks.max` 后仍然只有一个 Task？

这是该 Connector 的实现行为。它的 `taskConfigs` 始终返回一个 Task 配置，`tasks.max` 只提供 Kafka Connect 的任务上限，不会改变单实例的任务数量。需要隔离数据流或扩展部署时，为不同 Topic 集合创建多个 Connector 实例，并确保它们不会重复订阅同一 Topic。

### 为什么 `topics` 和 `topics.regex` 配置后无法启动？

Kafka Connect Sink 要求两者恰好有一个非空。删除其中一个，或将 `topics.regex` 改为合法且非空的 Java 正则表达式；如果同时配置了 DLQ Topic，还要确认该 Topic 不会被订阅列表或正则匹配。

### 为什么 Diffusion 中出现了预期之外的 Topic 路径？

检查 `diffusion.destination` 的令牌是否与记录实际拥有的键和 Schema 版本一致。`${topic}` 使用 Kafka 输入 Topic，`${key}` 使用非空记录键；记录键为空时版本令牌不会被解析，缺少可用版本时令牌也可能保留为字面量。应先使用固定路径或仅使用 `${topic}` 核对输入范围，再逐步加入键和版本令牌。

### 为什么发布成功后重启仍会看到同一条记录再次处理？

Diffusion 更新成功与 Kafka Offset 提交之间存在窗口。若连接在更新完成后、Offset 提交前中断，Kafka Connect 恢复时可能重新投递记录。检查目标路径是否使用了相同的键和路径，评估业务是否能接受重复处理，不要把 Diffusion 的最后写入结果视为应用层幂等保证。

### 为什么错误没有进入 DLQ？

DLQ 是 Kafka Connect 的通用错误处理能力，且只有配置 `errors.deadletterqueue.topic.name` 并满足相应错误处理条件时才会产生 DLQ 记录。Diffusion 发布失败由 Sink Task 的发布 Future 和刷新等待处理，仅配置 `errors.tolerance` 并不能保证这类失败自动写入 DLQ。应先查看 Connector、Task、错误和重试指标及日志，确认失败属于记录级转换错误还是目标端发布错误。
