Skip to main content

概述

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。以下示例消费 UTF-8 字符串消息,将其发布到固定 Channel。
替换 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
  • 重要级别:中
  • 有效值 / 注意事项truefalse;生产环境保持启用,避免以明文传输凭证和消息。

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 不会自动选择 80443

client.proxy.username

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

client.proxy.password

设置代理认证密码。
  • 类型PASSWORD
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:与代理用户名配套配置;通过安全配置机制注入,不记录真实值。

client.proxy.pref.auth.type

设置优先使用的代理认证类型。
  • 类型STRING
  • 默认值BASIC
  • 重要级别:中
  • 有效值 / 注意事项:仅允许 BASICDIGESTX_ABLY_TOKEN,区分大小写;选择代理实际支持的类型。

client.proxy.non.proxy.hosts

设置不经过代理的主机列表。
  • 类型LIST
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:多个主机用逗号分隔,由 SDK 的代理配置处理;不要假定支持任意通配符语法。

错误处理与 Push 选项

onFailedRecordMapping

设置 Channel 或消息名称映射失败时的处理策略。
  • 类型STRING
  • 默认值stop
  • 重要级别:中
  • 有效值 / 注意事项:仅允许 stopskipdlq,区分大小写。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 并行发布批次造成的先后交错,而不是追求跨实体的全局顺序。 配置示例:在快速开始配置中添加以下配置,覆盖默认的批执行线程数。
关键说明:该设置让一个 Task 串行执行批发布,但多个 Task、多个分区仍可能并发向同一 Channel 发布。它不是端到端严格有序或无重复的保障。将同一实体的事件写入同一 Kafka 分区,并在消息内容中携带实体版本或业务序号,由订阅端识别重复和过期状态;同时评估串行发布带来的吞吐变化。不要仅靠增大缓冲记录数应对持续积压,应结合消费延迟、请求耗时和 Worker 内存调整生产速率与发布并行度。

监控

监控内容

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

导入 Grafana 大盘

下载 Kafka Connect Grafana 大盘,确保 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 去重,保留可对账的事件来源;不要将单次正常恢复或串行发布配置理解为无条件不丢不重保障。