概述
Diffusion Source Connector 订阅 Diffusion 中符合 selector 的 Topic 更新,将每次收到的值转换为 Kafka Connect SourceRecord,并写入 Kafka。它位于 Diffusion 与 Kafka 之间,适合把实时状态、设备数据或应用事件引入 Kafka,供流处理、分析和下游服务消费。 每条 Kafka 记录的键是完整的 Diffusion Topic 路径,值来自 Diffusion JSON 数据。目标 Kafka Topic 可固定为一个名称,也可使用${topic} 根据 Diffusion Topic 路径动态生成;动态映射时路径中的 / 会替换为 _。Connector 采用单 Task 订阅和传输数据,不将 Diffusion Topic 自动拆分到多个 Task。
前置条件
- 要传输数据,Diffusion 中需要有与 selector 匹配的 JSON Topic;连接账号需要具备建立会话、选择并订阅这些 Topic 的权限。
- 目标 Kafka Topic 应满足名称规则并允许 Connector 写入;使用动态映射时,应提前检查路径转换后的 Topic 名称及可能的名称碰撞。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka、Diffusion 服务和需要订阅的 JSON Topic,确认网络连通和访问权限;创建与管理操作参见 AutoMQ 的管理 Connector。以下配置把匹配source/app/ 路径的更新汇聚到一个 Kafka Topic。
diffusion-events,记录键仍保留原始 Diffusion Topic 路径。密码应通过受控的配置或密钥管理机制提供,不要写入源码、日志或共享材料。JsonConverter 是否在 JSON 中包含 Schema 信息,取决于运行环境中该 Converter 的配置。
配置
Diffusion 连接与认证
diffusion.url
指定建立 Diffusion 会话使用的完整连接 URL。
- 类别:Diffusion 连接
- 类型:
string - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:填写 Diffusion 客户端可接受且 Worker 能访问的 URL,例如 WebSocket 地址。配置解析不校验 URL 语法、地址可达性或空字符串,实际连接失败会使 Task 启动失败。
diffusion.username
指定连接 Diffusion 时使用的 principal 用户名。
- 类别:Diffusion 认证
- 类型:
string - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:账号需要具备建立会话以及选择、订阅
diffusion.selector所匹配 Topic 的权限。配置解析不校验账号是否存在或是否具备权限。
diffusion.password
指定 Diffusion principal 的密码。
- 类别:Diffusion 认证
- 类型:
password - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:通过受控的密钥管理机制提供真实密码,避免出现在源码、日志或共享材料中。密码是否有效由 Diffusion 在连接时检查。
订阅范围与 Topic 映射
diffusion.selector
指定订阅 Diffusion Topic 的 selector。
- 类别:源数据选择
- 类型:
string - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:使用 Diffusion 支持的 selector 语法,并只匹配计划接入的 JSON Topic。selector 原样交给 Diffusion 解析;配置校验不检查语法、匹配结果或权限。订阅在启动阶段未能及时完成时,Task 会启动失败。
kafka.topic
指定 Diffusion Topic 路径到 Kafka Topic 名称的映射模式。
- 类别:目标 Topic
- 类型:
string - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:固定名称会把所有匹配更新汇聚到同一 Kafka Topic;
${topic}会替换为完整 Diffusion Topic 路径,随后结果中的/全部替换为_。其他 token 不会替换。最终名称必须符合 Kafka Topic 规则并具备写权限;不同路径可能映射为同名 Topic。
轮询与批量
diffusion.poll.interval
指定内存队列暂时没有记录时,每次再次检查前的等待时间,单位毫秒。
- 类别:轮询
- 类型:
int - 默认值:
1000 - 重要级别:低
- 有效值 / 注意事项:使用
0至2147483647。一次空轮询最多执行三次等待,因此返回空结果前的空闲等待可能接近该值的三倍。0仅取消检查之间的等待,不改变固定检查次数;负值会在运行时失败。
diffusion.poll.size
指定一次 poll 最多从内存队列返回的记录数。
- 类别:批量
- 类型:
int - 默认值:
128 - 重要级别:低
- 有效值 / 注意事项:使用
1至2147483647。该值只限制单次返回数量,不限制队列总大小,也不是 Kafka Producer 的batch.size。较大值可能增加单批转换和内存处理量;0或负值不能形成可用的批量行为。
运行与序列化
connector.class
指定 Diffusion Source Connector 实现类。
- 类别:Kafka Connect 框架配置
- 类型:
string - 默认值:无,必填
- 重要级别:高
- 有效值 / 注意事项:使用
com.diffusiondata.connect.diffusion.source.DiffusionSourceConnector。
tasks.max
指定 Connector 可使用的最大 Task 数。
- 类别:Kafka Connect 框架配置
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1。此 Connector 始终只生成一个 Task,提高该值不会拆分 selector、增加 Source 并行度或改变单 Task 行为,应保留为1。
tasks.max.enforce
指定 Kafka Connect 是否强制 Connector 生成的 Task 数不超过 tasks.max。
- 类别:Kafka Connect 框架配置
- 类型:
boolean - 默认值:
true - 重要级别:低
- 已弃用:是
- 替代项:无
- 有效值 / 注意事项:可取
true或false。Kafka Connect 已弃用该配置,应保留默认值;关闭限制不会使此 Connector 生成多个 Task。
key.converter
指定 Diffusion Topic 路径记录键的序列化 Converter。
- 类别:Kafka Connect 框架配置
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略时继承 Worker 的 Key Converter。记录键为字符串形式的完整 Diffusion Topic 路径;需要按纯文本序列化记录键时可使用
org.apache.kafka.connect.storage.StringConverter,该类不是默认值。
value.converter
指定 Diffusion JSON 转换后的记录值使用的序列化 Converter。
- 类别:Kafka Connect 框架配置
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:省略时继承 Worker 的 Value Converter。Connector 会先把 JSON 值转换为 Kafka Connect 数据及动态推断的 Schema;需要 JSON 输出时可使用
org.apache.kafka.connect.json.JsonConverter,该类不是默认值。Converter 专属设置由所选 Converter 和运行环境管理。
最佳实践
按 Diffusion 路径拆分目标 Kafka Topic
适用业务场景:已使用固定 Kafka Topic 完成初始连通性验证,需要让不同 Diffusion Topic 的数据进入不同 Kafka Topic,以便下游按业务流分别设置保留策略、权限或消费任务,同时保留来源路径与目标名称之间的直观对应关系。 配置示例:以快速开始为基础,将kafka.topic 改为以下模式,其余配置保持不变。
source/app/orders 会映射为 diffusion_source_app_orders。上线前列出 selector 可能匹配的路径,检查转换后的名称是否合法且互不冲突,并准备相应 Kafka Topic 或自动创建策略与写权限。映射只替换 ${topic} 和 /,不会清理其他非法字符;foo/bar 与 foo_bar 可能得到同一结果。若业务需要把所有更新汇聚到同一 Topic,应继续使用快速开始中的固定名称,不要同时追求动态拆分。
运行期间调整等待时间与单批数量
适用业务场景:Connector 已稳定传输数据,需要在空闲期新数据等待时间、单次处理量和 Worker 内存压力之间进行调整。 配置示例:在现有 Connector 配置中覆盖以下两项,其余连接、selector、Topic 映射和 Converter 配置保持不变。监控
监控内容
关注 Kafka Connect 集群健康、Connector 与 Task 状态、记录吞吐、端到端延迟、Offset 提交、错误和重试,以及 Worker JVM 内存、GC 和线程信号;Task 处于 RUNNING 不代表 selector 持续收到更新,应结合 Diffusion 源数据变化和 Kafka 实际输出判断。由于 Connector 内部队列没有容量上限,还应特别关注持续流量下的堆内存增长。只有部署启用了相应错误处理时才关注 DLQ 活动。导入 Grafana 大盘
下载共享的 AutoMQ Connect Cluster Dashboard,确认 Grafana 数据源能够查询 Kafka Connect 和 Worker JVM 指标,并与大盘所需的集群、Connector、Task 等标签匹配,再通过 Grafana 的导入功能加载 JSON 并选择对应数据源。限制条件
- 单个 Connector 实例始终只有一个 Task,
tasks.max不能拆分 selector 或提高并行度。 - SourceRecord 不包含可恢复的 source offset,Connector 不能从 Kafka Connect checkpoint 恢复 Diffusion 更新位置,也不能回放 Task 停止期间的更新。
- Connector 不提供可据此声明的 exactly-once、at-least-once 或无重复交付保障;业务应根据可重放性和下游幂等要求评估使用方式。
- Diffusion 回调与 poll 之间使用无界内存队列,没有源端背压或总积压上限;持续输入快于 Kafka 发送时,内存占用可能持续增长。
- 每条 JSON 消息独立推断值 Schema,数据形状变化可能改变或移除 Schema;JSON 对象不会恢复为预定义的 Connect Struct。
- Topic 映射只将
/替换为_,不会清理其他非法字符或检测名称碰撞。 - Connector 创建 Diffusion 会话时不启用客户端自动重连,也没有内建的重新连接和重新订阅循环。
常见问题
Task 正常运行,但 Kafka Topic 中没有新消息怎么办?
先确认 Diffusion 中确实有新的 JSON Topic 值更新,并检查diffusion.selector 是否匹配对应路径。再核对 principal 的选择与订阅权限、kafka.topic 的最终名称、Kafka 写权限以及 Key 和 Value Converter。selector 没有匹配项不会在配置校验阶段报错;使用固定 Topic 时不要误查动态派生名称,使用动态映射时应按路径转换结果查找。查看 Worker 日志中的订阅、转换和 Kafka 写入错误,但不要输出真实密码。
为什么不同 Diffusion Topic 的数据进入了同一个 Kafka Topic?
如果kafka.topic 是固定名称,所有 selector 匹配的数据都会汇聚到该 Topic。需要按来源拆分时使用包含 ${topic} 的模式,并提前检查转换结果。即使使用动态模式,foo/bar 与 foo_bar 也会因 / 替换为 _ 而发生名称碰撞;调整源路径规划或添加不会冲突的固定前缀,并逐一核对最终 Topic 名称。