Skip to main content

概述

JR Source Connector 调用 Kafka Connect Worker 上安装的 JR 命令行工具,按配置生成随机 JSON 对象,并将结果作为 Kafka Connect SourceRecord 写入一个 Kafka Topic。它适合为开发、测试和集成环境持续提供可控的模拟数据,也支持从自定义模板生成符合目标数据结构的记录。 Connector 每次运行可以生成一个或多个对象,并按设定的最小执行间隔周期性生成数据。默认情况下消息只有 Value;也可以从模板字段或单独的键模板生成消息 Key。生成结果支持 String、Avro、JSON Schema 和 Protobuf Converter。

前置条件

  • JR 必须安装在所有可能运行该 Task 的 Kafka Connect Worker 上,并确保 Worker 能够执行 jr;如果未将其加入 PATH,请配置 jr_executable_path

授权许可

使用 MIT License。

快速开始

请先准备 Connect Cluster、Kafka Topic 和已安装 JR 的 Worker,确认 Worker 能访问目标 Topic。Connector 的创建和管理方式请参阅 AutoMQ 的管理 Connector 下面的配置使用内置的 net_device 模板,将生成的数据写入 jr-generated-data Topic:
topic 替换为实际目标 Topic。net_device 必须是 Worker 上 JR 可用的内置模板;如果使用自定义模板,请改用 embedded_template。如果 JR 不在 Worker 的 PATH 中,还需要设置 jr_executable_path

配置

模板

template

选择 JR 内置的 Value 模板名称。
  • 类型string
  • 默认值net_device
  • 重要级别:高
  • 有效值 / 注意事项:当 embedded_template 为空时,Connector 启动会执行 jr list -n,并要求该名称出现在返回的模板列表中。embedded_template 非空时优先使用自定义模板,跳过该检查。

embedded_template

从 Worker 可访问的文件或 HTTP(S) URL 加载自定义 Value 模板。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:非空值优先于 template。除 HTTP(S) URL 外的值按文件路径处理;文件和 URL 必须能被运行 Task 的每个 Worker 访问。读取内容中的换行会在传给 JR 前移除。

输出

topic

指定生成记录写入的 Kafka Topic。
  • 类型list
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须解析为且只能包含一个非空 Topic。该 Connector 不支持通过此配置路由到多个 Topic。

调度与生成

frequency

设置两次生成检查之间的最小间隔,单位为毫秒。
  • 类型long
  • 默认值5000
  • 重要级别:高
  • 有效值 / 注意事项:建议使用正数。Task 只有在当前时间严格晚于上一次开始执行 JR 的时间加上该值时才生成下一批数据;因此这是两次 JR 命令开始时间之间的最小间隔,不会在命令执行完成后再额外等待完整间隔。配置解析不会拒绝非正数,非正数可能导致频繁轮询。

duration

限制 Task 从启动开始允许生成数据的总时长,单位为毫秒。
  • 类型long
  • 默认值-1
  • 重要级别:中
  • 有效值 / 注意事项:大于 1 的值建立截止时间;计时从 Task start() 开始,第一次符合 frequency 的生成仍会执行。省略该配置或将其设置为不大于 1 时,不会建立截止时间;其中小于 1 的值会被归一化为 -1。因此 duration=1 不表示运行 1 毫秒,而是无限运行。

objects

设置每次调用 JR 生成的对象数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:小于 1 的值会被归一化为 1。增大该值会增加单次批量的内存占用、JSON 解析和序列化工作。

消息键

key_field_name

使用 Value 模板中的字段生成消息 Key。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:非空时,Connector 请求 JR 交替输出 Key JSON 和 Value JSON,并使用指定字段生成键值替换。实现不会在启动时检查该字段是否存在;如果未配置该项且未配置 key_embedded_template,消息没有 Key。

key_value_interval_max

设置由 key_field_name 生成的整数键值上限。
  • 类型int
  • 默认值100
  • 重要级别:中
  • 有效值 / 注意事项:仅在未配置 key_embedded_template 且使用 key_field_name 时生效。小于 1 的值会被归一化为 100;生成值范围为 0 到该上限。

key_embedded_template

从 Worker 可访问的文件或 HTTP(S) URL 加载自定义 Key 模板。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:非空时优先于 key_field_namekey_value_interval_max。加载规则与 embedded_template 相同,要求 JR 输出交替的 Key JSON 和 Value JSON。

运行时

jr_executable_path

指定 jr 可执行文件所在的目录。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:配置的是目录,不是包含文件名的完整路径。Connector 会在该目录后追加路径分隔符和 jr;省略时使用 Worker 的 PATH。每个可能运行 Task 的 Worker 都必须安装并允许执行 JR。

序列化

value.converter

指定生成记录 Value 使用的 Kafka Connect Converter。
  • 类型string
  • 默认值org.apache.kafka.connect.storage.StringConverter
  • 重要级别:中
  • 有效值 / 注意事项:支持 StringConverterAvroConverterJsonSchemaConverterProtobufConverter。后三者需要 Worker 中存在对应依赖,并按 Converter 要求配置 Schema Registry;相关 Schema Registry 配置不属于本 Connector 的配置项。

key.converter

指定生成记录 Key 使用的 Kafka Connect Converter。
  • 类型string
  • 默认值org.apache.kafka.connect.storage.StringConverter
  • 重要级别:中
  • 有效值 / 注意事项:仅在生成消息 Key 时有实际作用。支持 StringConverterAvroConverterJsonSchemaConverterProtobufConverter;使用 Schema Converter 时,需要在 Worker 中提供匹配的依赖和 Schema Registry 配置。

Kafka Connect 框架

connector.class

指定 Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须设置为 io.jrnd.kafka.connect.connector.JRSourceConnector

tasks.max

设置 Kafka Connect 为 Connector 请求的最大 Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:该 Connector 始终返回一个 Task 配置,且只支持一个目标 Topic;增大 tasks.max 不会创建多个 JR 生成任务,也不会自动进行 Source 数据分片。

最佳实践

在限定时间内生成一批测试数据

适用业务场景:需要为一次集成测试、演示或阶段性数据准备生成有限时长的数据,而不是让 Source 持续运行。配置截止时间后,Task 从启动开始计时,在第一次满足 frequency 的轮询中仍会生成数据,之后到达截止时间就不再继续生成。 配置示例
topic 替换为测试使用的 Topic。duration 大于 1 时才启用有限时长;它只停止继续生成数据,不代表 Connector 或 Task 自动停止。 关键说明duration 的计时起点是 Task start(),不是第一次生成。若 frequency 大于剩余时长,第一次符合条件的轮询仍可能先执行一次;需要更短的准备窗口时,应同时调整 frequencyduration,不要使用 duration=1 作为短时运行配置。

使用自定义模板匹配下游数据结构

适用业务场景:内置模板不能提供下游测试或集成所需的字段结构,需要从 Worker 本地文件或可访问的 HTTP(S) 地址加载自定义 JR 模板,并将生成的 JSON 写入一个 Topic。 配置示例
embedded_template 替换为每个 Worker 都能访问的文件路径或 HTTP(S) URL。配置自定义模板后不需要再设置 templateembedded_template 非空时会覆盖它。 关键说明:Connector 会移除模板内容中的换行后再传给 JR,因此模板应保持可被 JR 解析的 JSON 模板格式。文件或 URL 读取发生在 Connector 启动时,部署到多个 Worker 时要确保路径、网络访问和文件内容一致。

监控

监控内容

监控 Kafka Connect Worker、Connector 和 Task 的健康状态与状态转换,关注生成记录吞吐、轮询到记录的延迟、Source offset 提交进度、错误和重试次数,以及 Worker JVM 的堆内存、垃圾回收、CPU 和线程信号。若部署启用了错误容忍或死信队列,再关注对应的错误处理和 DLQ 活动;同时检查 JR 子进程异常是否导致 Task 进入失败状态。

导入 Grafana 大盘

下载 AutoMQ Connect Cluster Dashboard,在 Grafana 中使用已采集 Kafka Connect 和 Worker 指标的数据源导入该 JSON,并按实际集群标签选择 Worker、Connector 和 Task。

限制条件

  • 只能配置一个非空目标 Topic,不支持由 Connector 将记录路由到多个 Topic。
  • tasks.max 不会把生成工作拆分为多个 Task;Connector 实际返回一个 Task 配置。
  • JR 必须安装并可执行于所有可能运行 Task 的 Worker;jr_executable_path 配置的是目录而不是完整可执行文件路径。
  • embedded_templatekey_embedded_template 依赖 Worker 对文件或 HTTP(S) 资源的访问,不能把只在提交配置的客户端上可访问的路径当作 Worker 本地路径。
  • duration 只有大于 1 的值才限制运行时长;省略该配置或将其设置为不大于 1 时会无限运行,且有限时长从 Task 启动开始计时。
  • Source offset 的 position 不是 JR 记录 ID、时间戳或内容哈希。实现按完整的高精度时间字符串更新内部计数,并在创建记录时再次递增,因此该值可能重置或重复,常见值为 2;不能将其当作全局递增序号。
  • 生成任务每次重新调用 JR 获取随机数据,不会恢复此前的随机序列,也不提供 Connector 级别的去重、重试或死信队列策略。

常见问题

Connector 启动时报模板不存在,如何处理?

embedded_template 为空时,Connector 会调用 jr list -n 校验 template。请在运行 Task 的 Worker 上执行同一 JR 安装,确认模板名称存在且 JR 可执行;或者改用 Worker 可访问的文件或 HTTP(S) 地址配置 embedded_template。如果设置了非空 embedded_template,它会优先于 template

配置了 duration,为什么没有立即停止?

只有大于 1duration 才会建立截止时间,并且计时从 Task 启动开始。duration=1 不会建立截止时间,因此表现为无限运行。有限时长到达后,Connector 停止继续生成数据,但不会自动删除或停止 Connector、Task;应通过 Connect 管理操作处理其生命周期。

为什么 Source offset 没有按记录数持续递增?

该 Connector 的 position 不是业务记录序号。实现根据完整的 UTC ISO 时间字符串更新内部计数:时间变化时先将计数重置为 1,创建 SourceRecord 时还会再次递增,因此连续记录可能重复得到 position=2,而不是持续单调递增。不要使用该字段作为业务唯一标识;需要稳定关联时,应配置消息 Key,并由下游按业务字段处理。

为什么生成了记录但消息没有 Key?

默认配置只生成 Value。要生成 Key,请配置 key_field_name,或使用 key_embedded_template 提供独立的 Key 模板;如果使用后者,它会覆盖 key_field_name。同时确认 Key Converter 与生成的 Key 结构匹配。

为什么增大 tasks.max 后没有出现多个生成任务?

该 Connector 的任务配置固定返回一个 Task,且输出目标限制为一个 Topic。tasks.max 是 Kafka Connect 的请求上限,不会改变这个 Connector 的任务分片逻辑;需要调整生成量时,应结合 frequencyobjects 评估单个 Task 的批量与执行节奏。