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

# Ably Sink Connector

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

## 概述

Ably Sink Connector 将 Kafka Topic 中的事件发布到 Ably Channel，连接后端事件流与实时订阅客户端，适用于状态更新、业务事件通知和实时看板。每条 Kafka 记录对应一条 Ably 消息，目标 Channel 可以固定，也可以根据 Topic、分区或记录字段动态选择；消息名称可用于区分事件类型。

消息内容取决于 Kafka Connect Converter 的输出。带 Schema 的 Struct 会转换为 JSON 内容，不附带 Schema 定义；字符串和字节内容走相应的消息转换路径。无 Schema 的 Map 或 List 不会自动转换为结构化 JSON，因此需要 JSON 内容时，应提供 JSON 字符串或带 Schema 的 Struct。

恢复消费依赖 Kafka Connect 已提交的 Offset。发布成功但尚未提交 Offset 时发生故障，或请求响应丢失后重试，都可能产生重复消息。不能将此 Connector 视为恰好一次投递或无条件至少一次投递方案；对于不能遗漏或重复处理的业务，应在事件内容中保留业务事件 ID，并配合下游去重、对账和补发机制。

## 前置条件

* Ably API Key 对所有静态或动态目标 Channel 具备 `publish` 权限。
* 使用嵌套字段模板时，Converter 必须提供对应的 Connect Schema 和 Struct。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka Topic 和 Ably 账号及目标 Channel，并确认网络连通与访问权限。集群准备和 Connector 管理操作参见[管理 Connector](../manage-connectors)。以下示例消费 UTF-8 字符串消息，将其发布到固定 Channel。

```properties theme={null}
connector.class=com.ably.kafka.connect.ChannelSinkConnector
topics=<kafka-topic>
channel=<ably-channel>
client.key=<ably-api-key>
client.id=<ably-client-id>
value.converter=org.apache.kafka.connect.storage.StringConverter
```

替换 Topic、Channel、API Key 和客户端标识后，将上述内容作为 Connector 配置应用。`client.id` 是非空的 Ably 客户端身份，不是 Kafka 消费者客户端标识。示例不依赖记录 Key，因此沿用 Worker 的 Key Converter；若要使用 Key 模板，需另外确认 Key 解码方式。凭证应通过部署环境的安全配置机制注入，不要写入版本库或共享日志。

## 配置

下面列出全部公开插件配置，以及直接影响订阅、并行度和消息解码的六项 Kafka Connect 框架配置。默认值“无”表示没有默认值，必填项需显式配置；`null` 表示未指定，`[]` 表示空列表，二者不等同于空字符串。配置名称区分大小写。

### Connector 与订阅

#### `connector.class`

指定 Sink Connector 实现类。

* **类型**：`STRING`
* **默认值**：无
* **必填**：是
* **重要级别**：高
* **有效值 / 注意事项**：使用 `com.ably.kafka.connect.ChannelSinkConnector`，所有运行 Task 的 Worker 均需能够加载该插件。

#### `topics`

指定需要消费的 Kafka Topic。

* **类型**：`LIST`
* **默认值**：`[]`
* **重要级别**：高
* **有效值 / 注意事项**：多个名称用逗号分隔；与 `topics.regex` 二选一且必须配置其中一项。启用 DLQ 时不能订阅 DLQ Topic。

#### `topics.regex`

通过正则表达式选择需要消费的 Kafka Topic。

* **类型**：`STRING`
* **默认值**：`""`
* **重要级别**：高
* **有效值 / 注意事项**：使用 Java 正则语法；与 `topics` 互斥，启用 DLQ 时不能匹配 DLQ Topic。

#### `tasks.max`

设置最大 Task 数量。

* **类型**：`INT`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：至少为 `1`；实际消费并行度受订阅分区数约束。同一分区不会同时分配给同一消费组的多个 Task，与每个 Task 的批执行线程数分别控制。

### 认证与日志

#### `client.key`

设置用于发布消息的 Ably API Key。

* **类型**：`PASSWORD`
* **默认值**：无
* **必填**：是
* **重要级别**：高
* **有效值 / 注意事项**：使用有效 API Key，权限需覆盖所有目标 Channel；不要将真实值暴露在配置示例、版本库或日志中。

#### `client.id`

设置 Ably 客户端身份。

* **类型**：`STRING`
* **默认值**：无
* **必填**：是
* **重要级别**：高
* **有效值 / 注意事项**：非空字符串；不是 Kafka 消费者的 `client.id`。

#### `client.loglevel`

设置 Ably SDK 日志级别。

* **类型**：`INT`
* **默认值**：`2`
* **重要级别**：低
* **有效值 / 注意事项**：已定义级别为 `2`（VERBOSE）、`3`（DEBUG）、`4`（INFO）、`5`（WARN）、`6`（ERROR）、`99`（NONE）。生产环境按排障需要选择，避免长期输出过量调试信息。

### 消息解码与映射

#### `key.converter`

设置 Kafka 记录 Key 的 Converter。

* **类型**：`CLASS`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未指定时继承 Worker 设置。使用 `#{key}` 或 `#{key.<path>}` 时，需选择与实际编码一致的 Converter；嵌套路径要求 Schema/Struct。

#### `value.converter`

设置 Kafka 记录 Value 的 Converter。

* **类型**：`CLASS`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未指定时继承 Worker 设置。字符串可使用 `org.apache.kafka.connect.storage.StringConverter`；`#{value.<path>}` 需要 Schema/Struct，不能仅提供 JSON 文本或无 Schema Map。Converter 自身的选项需按所选 Converter 配置。

#### `channel`

设置目标 Ably Channel 的静态名称或模板。

* **类型**：`STRING`
* **默认值**：无
* **必填**：是
* **重要级别**：高
* **有效值 / 注意事项**：非空字符串，首字符不能为冒号、逗号、空白或 `[`，名称不能包含换行。模板支持 `#{topic}`、`#{topic.name}`、`#{topic.partition}`、`#{key}`、`#{key.<path>}` 和 `#{value.<path>}`。嵌套路径沿 Struct 字段读取，不支持数组索引或 Map 路径；字段缺失等映射错误受 `onFailedRecordMapping` 控制。模板输出也应满足 Ably Channel 命名及权限要求。

#### `message.name`

设置 Ably 消息名称，用于标识事件类型。

* **类型**：`STRING`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：可使用静态文本或与 `channel` 相同的模板；未指定或为空时不设置名称。没有直接 `#{value}` 占位符。

#### `messagePayloadSizeMax`

消息载荷大小相关的公开参数。

* **类型**：`INT`
* **默认值**：`65536`
* **重要级别**：中
* **有效值 / 注意事项**：此项不执行载荷大小限制，不应依赖它校验、截断或拆分消息。应在上游控制实际消息和请求大小，并遵守 Ably 服务限制。

### 批处理与并发

#### `batchExecutionThreadPoolSize`

设置每个 Task 的批发布执行线程数。

* **类型**：`INT`
* **默认值**：`10`
* **重要级别**：中
* **有效值 / 注意事项**：必须大于 `0`。设置为 `1` 可串行执行该 Task 的批发布；多个线程可使批次并行完成，不能保证跨批完成顺序，也不能提供跨 Task 全局顺序。

#### `batchExecutionMaxBufferSize`

设置提交一个批次前缓冲的最大 Kafka 记录数。

* **类型**：`INT`
* **默认值**：`100`
* **重要级别**：中
* **有效值 / 注意事项**：使用大于 `0` 的值；达到数量阈值即提交批次。此项不是字节大小上限，也不是所有待发送数据的积压上限。

#### `batchExecutionMaxBufferSizeMs`

设置未满批缓冲的定时发送等待时间，单位毫秒。

* **类型**：`INT`
* **默认值**：`100`
* **重要级别**：中
* **有效值 / 注意事项**：使用非负值；`0` 可立即触发定时发送。满批可提前发送，此项不是每条记录的最低延迟，也不是端到端延迟上限。

#### `client.async.http.threadpool.size`

设置 Ably SDK 异步 HTTP 线程池大小。

* **类型**：`INT`
* **默认值**：`64`
* **重要级别**：中
* **有效值 / 注意事项**：这是 SDK 异步请求设置，不是此 Connector 的批发布并发控制；批发布并发使用 `batchExecutionThreadPoolSize` 调整。

### 连接与备用端点

#### `client.tls`

设置是否使用 TLS 连接 Ably 服务。

* **类型**：`BOOLEAN`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：`true` 或 `false`；生产环境保持启用，避免以明文传输凭证和消息。

#### `client.rest.host`

设置自定义 REST 主机。

* **类型**：`STRING`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未指定时由 SDK 选择 REST 主机；仅在确有自定义端点需求时设置，并确认认证、网络和备用端点配置一致。使用非默认 REST 主机时，不要同时设置 `client.environment`。

#### `client.port`

设置非 TLS 连接的端口。

* **类型**：`INT`
* **默认值**：`0`
* **重要级别**：低
* **有效值 / 注意事项**：`0` 表示使用 SDK 的非 TLS 默认端口 `80`；仅在 `client.tls=false` 时使用。显式设置时应使用有效服务端口。

#### `client.tls.port`

设置 TLS 连接的端口。

* **类型**：`INT`
* **默认值**：`0`
* **重要级别**：低
* **有效值 / 注意事项**：`0` 表示使用 SDK 的 TLS 默认端口 `443`；仅在 `client.tls=true` 时使用。显式设置时应使用有效服务端口。

#### `client.environment`

设置非默认的 Ably 服务环境。

* **类型**：`STRING`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：仅在 Ably 服务部署要求时配置，端点选择交由 SDK 处理；不能与非默认 `client.rest.host` 同时使用。它不是 Kafka Connect 或托管平台的环境标识。

#### `client.http.open.timeout`

设置建立 HTTP 连接的超时时间，单位毫秒。

* **类型**：`INT`
* **默认值**：`4000`
* **重要级别**：中
* **有效值 / 注意事项**：按网络条件设置合理的非负超时，不要将此项作为整个消息处理流程的超时。

#### `client.http.request.timeout`

设置单次 HTTP 请求及响应的超时时间，单位毫秒。

* **类型**：`INT`
* **默认值**：`10000`
* **重要级别**：中
* **有效值 / 注意事项**：按网络和请求规模设置合理的非负超时；不包含整个 Kafka Connect 错误重试窗口。

#### `client.http.max.retry.count`

设置 SDK HTTP 备用主机重试次数。

* **类型**：`INT`
* **默认值**：`3`
* **重要级别**：中
* **有效值 / 注意事项**：建议使用非负值；重试需要可用备用主机且错误符合 SDK 主机失败条件。此项不会让所有 HTTP 错误重试，也不是 Kafka 消费重试策略。

#### `client.fallback.hosts`

设置 HTTP 请求可使用的备用主机列表。

* **类型**：`LIST`
* **默认值**：`[]`
* **重要级别**：中
* **有效值 / 注意事项**：多个主机用逗号分隔，必须与实际 Ably 环境匹配。空列表表示没有备用候选，不会自动启用 SDK 内置备用主机；仅增加重试次数不足以启用备用主机重试。

### 代理

#### `client.proxy`

设置是否启用 HTTP 代理。

* **类型**：`BOOLEAN`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：启用时必须提供代理主机和实际端口。未启用代理时，代理参数的类型和枚举值仍需合法。

#### `client.proxy.host`

设置 HTTP 代理主机。

* **类型**：`STRING`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：`client.proxy=true` 时需设置非空、可达的代理主机。

#### `client.proxy.port`

设置 HTTP 代理端口。

* **类型**：`INT`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：启用代理时必须设置实际有效端口；默认 `0` 不会自动选择 `80` 或 `443`。

#### `client.proxy.username`

设置代理认证用户名。

* **类型**：`STRING`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：不是 Ably 身份；启用代理并设置用户名时，还必须设置 `client.proxy.password`。

#### `client.proxy.password`

设置代理认证密码。

* **类型**：`PASSWORD`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：与代理用户名配套配置；通过安全配置机制注入，不记录真实值。

#### `client.proxy.pref.auth.type`

设置优先使用的代理认证类型。

* **类型**：`STRING`
* **默认值**：`BASIC`
* **重要级别**：中
* **有效值 / 注意事项**：仅允许 `BASIC`、`DIGEST`、`X_ABLY_TOKEN`，区分大小写；选择代理实际支持的类型。

#### `client.proxy.non.proxy.hosts`

设置不经过代理的主机列表。

* **类型**：`LIST`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：多个主机用逗号分隔，由 SDK 的代理配置处理；不要假定支持任意通配符语法。

### 错误处理与 Push 选项

#### `onFailedRecordMapping`

设置 Channel 或消息名称映射失败时的处理策略。

* **类型**：`STRING`
* **默认值**：`stop`
* **重要级别**：中
* **有效值 / 注意事项**：仅允许 `stop`、`skip`、`dlq`，区分大小写。`stop` 停止批处理，`skip` 丢弃映射失败记录；`dlq` 需要可用的 Kafka Connect 错误记录报告机制、已配置的 `errors.deadletterqueue.topic.name` 和允许持续处理的 `errors.tolerance=all`。缺少报告机制时不能继续按 DLQ 模式处理。此项不覆盖全部 Converter、内容转换或 HTTP 发布错误。

#### `client.push.full.wait`

设置传递给 Ably SDK 的 Push REST 操作等待选项。

* **类型**：`BOOLEAN`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：此 Sink 不执行设备 Push 管理操作，该项不控制批发布完成，也不保证设备通知送达。需要通知业务时，还须独立完成 Ably Push 配置和设备注册。

## 最佳实践

### 降低状态事件的批次交错

**适用业务场景**：订单或设备状态事件已按业务实体分区，希望减少同一 Task 并行发布批次造成的先后交错，而不是追求跨实体的全局顺序。

**配置示例**：在快速开始配置中添加以下配置，覆盖默认的批执行线程数。

```properties theme={null}
batchExecutionThreadPoolSize=1
```

**关键说明**：该设置让一个 Task 串行执行批发布，但多个 Task、多个分区仍可能并发向同一 Channel 发布。它不是端到端严格有序或无重复的保障。将同一实体的事件写入同一 Kafka 分区，并在消息内容中携带实体版本或业务序号，由订阅端识别重复和过期状态；同时评估串行发布带来的吞吐变化。不要仅靠增大缓冲记录数应对持续积压，应结合消费延迟、请求耗时和 Worker 内存调整生产速率与发布并行度。

## 监控

### 监控内容

关注 Kafka Connect 集群和 Worker 健康状态、Connector / Task 状态、吞吐、消费与处理延迟、Offset 提交及失败、错误和重试，并结合 Worker JVM 堆内存、GC 与线程情况判断积压。Task 处于运行状态或 Offset 已提交不能单独证明所有消息已被业务客户端接收，应配合业务事件对账；仅在启用了相应错误处理时关注 DLQ 活动。

### 导入 Grafana 大盘

下载 [Kafka Connect Grafana 大盘](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确保 Prometheus 兼容数据源已采集 Kafka Connect 指标且集群、Worker 等标签与大盘查询一致，然后在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* 不提供 Kafka 到 Ably 的恰好一次投递保障，不应依赖自动幂等发布消除重复。
* 不提供跨分区、跨 Task 或跨 Channel 的全局顺序。
* 带 Schema 的 Value 仅支持顶层 Struct、String 和 Bytes，不支持顶层数字、布尔、Array 或 Map。
* 嵌套模板要求 Schema/Struct，不支持解析 JSON 字符串、无 Schema Map、数组索引或 Map 路径。
* 模板叶值支持字符串、Integer、Long、Boolean 和有效 UTF-8 的字节数组，不支持 Float、Double、Byte、Short 或 ByteBuffer。
* Struct 转 JSON 不支持直接的原始字节数组字段，不应假定任意逻辑类型或时间精度都能原样保留。
* 无 Value Schema 时，不导出 Kafka Key、Headers 或 Push Extras；带 Schema 时也不会自动导出 Topic、分区、Offset 或 Kafka 时间戳。
* 空 Value 不表示删除 Ably 数据，也不会自动跳过 Kafka tombstone。
* 缓冲记录数不能限制总待发送积压，持续慢请求可能导致内存压力。
* `messagePayloadSizeMax` 不限制实际载荷，消息大小、批请求大小及速率仍受 Ably 服务约束。

## 常见问题

### Task 运行中，为什么目标 Channel 没有消息？

检查 Topic 是否有新记录、消费组位置是否已到末尾、订阅是否匹配，以及实际模板输出的 Channel 是否与订阅端一致。确认 API Key 对该 Channel 有发布权限，再检查 REST 错误、连接超时和 Worker 内存。低流量时还应考虑未满批定时发送的等待时间；运行状态本身不代表发布成功。

### JSON 中有字段，为什么字段模板仍然失败？

JSON 字符串和无 Schema Map 不能作为嵌套模板的数据来源。检查 Converter 输出是否为带 Schema 的 Struct、字段路径是否存在，以及叶值是否为支持的类型；按实际 Kafka 编码调整 Converter 或上游事件结构。字段重命名和删除也会影响模板，发布 Schema 变更前应同步检查引用字段。

### 为什么消息内容不是预期的 JSON 对象，或者没有 Key 和 Headers？

无 Schema 的 Map/List 不会自动发布为结构化 JSON，无 Value Schema 时也不会生成 Key 和 Headers Extras。需要 JSON 内容时使用 JSON 字符串或带 Schema 的 Struct；需要 Extras 时使用受支持的带 Schema Value，并检查 Key 和 Header 的解码方式。字符串 Key 保持字符串，字节数组 Key 使用 Base64，其他 Key 类型不会自动导出。

### 已增加重试次数，为什么发布失败后仍没有重试？

检查 `client.fallback.hosts` 是否配置了与环境一致的可用备用主机；默认空列表没有备用候选。SDK 仅对符合主机失败条件的错误进行备用主机重试，不应期待权限错误、限流或所有服务端错误都自动重试。先处理凭证、权限或速率问题，不能通过增大重试次数替代这些修复。

### 为什么映射失败后没有记录进入 DLQ？

确认 `onFailedRecordMapping=dlq`，并检查 Kafka Connect 的 DLQ Topic、写权限、错误容忍设置以及错误记录报告机制是否可用；仅设置插件策略不足以启用完整 DLQ 处理。还应区分映射失败与 Converter、内容格式或 HTTP 发布错误，该策略不覆盖全部错误类型。不要使用 `skip` 掩盖必须保留的业务事件，应为 DLQ 配置告警、修复和重放流程。

### 重启后为什么收到重复事件？

发布与 Offset 提交不是同一个原子操作，已发布但尚未提交的记录可能再次消费，请求响应丢失后的重试也可能重复。使用稳定业务事件 ID 去重，保留可对账的事件来源；不要将单次正常恢复或串行发布配置理解为无条件不丢不重保障。
