概述
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 - 重要级别:中
- 有效值 / 注意事项:可选值为
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 列表。 配置示例:topics.regex 非空时不要同时配置非空的 topics。${topic} 会保留输入 Topic 名称,因此 orders-created 会写入 events/orders-created。订阅范围扩大前,应确认 Diffusion principal 对所有可能生成的目标路径都有更新或创建权限。
用多个 Connector 实例分隔数据流
适用业务场景:需要分开管理不同业务流,分别设置目标路径或故障边界,且单个 Connector 实例中的 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,在 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、错误和重试指标及日志,确认失败属于记录级转换错误还是目标端发布错误。