Skip to main content

概述

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
<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
  • 重要级别:中
  • 有效值 / 注意事项:可选值为 noneall。该配置是通用错误策略;仅配置它并不能保证 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 列表。 配置示例
关键说明topics.regex 非空时不要同时配置非空的 topics${topic} 会保留输入 Topic 名称,因此 orders-created 会写入 events/orders-created。订阅范围扩大前,应确认 Diffusion principal 对所有可能生成的目标路径都有更新或创建权限。

用多个 Connector 实例分隔数据流

适用业务场景:需要分开管理不同业务流,分别设置目标路径或故障边界,且单个 Connector 实例中的 Topic 集合已经较大。此时可为每组 Topic 创建独立的 Connector 实例。 配置示例
为另一组 Topic 创建第二个 Connector 实例时,应使用不同的 topicstopics.regex 和目标路径,而不是只提高同一实例的 tasks.max。该实现始终只创建一个 Task,增加 tasks.max 不会在实例内部产生任务级并行。多个实例之间应明确划分输入 Topic,避免重复消费同一数据流。

监控

监控内容

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

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • 单个 Diffusion Sink Connector 实例始终只创建一个 Task;提高 tasks.max 不会让该实例并行创建多个 Task,需要通过多个实例分隔输入 Topic 才能进行部署级扩展。
  • topicstopics.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。

为什么 topicstopics.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、错误和重试指标及日志,确认失败属于记录级转换错误还是目标端发布错误。