概述
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的值建立截止时间;计时从 Taskstart()开始,第一次符合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_name和key_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 - 重要级别:中
- 有效值 / 注意事项:支持
StringConverter、AvroConverter、JsonSchemaConverter和ProtobufConverter。后三者需要 Worker 中存在对应依赖,并按 Converter 要求配置 Schema Registry;相关 Schema Registry 配置不属于本 Connector 的配置项。
key.converter
指定生成记录 Key 使用的 Kafka Connect Converter。
- 类型:
string - 默认值:
org.apache.kafka.connect.storage.StringConverter - 重要级别:中
- 有效值 / 注意事项:仅在生成消息 Key 时有实际作用。支持
StringConverter、AvroConverter、JsonSchemaConverter和ProtobufConverter;使用 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 大于剩余时长,第一次符合条件的轮询仍可能先执行一次;需要更短的准备窗口时,应同时调整 frequency 和 duration,不要使用 duration=1 作为短时运行配置。
使用自定义模板匹配下游数据结构
适用业务场景:内置模板不能提供下游测试或集成所需的字段结构,需要从 Worker 本地文件或可访问的 HTTP(S) 地址加载自定义 JR 模板,并将生成的 JSON 写入一个 Topic。 配置示例:embedded_template 替换为每个 Worker 都能访问的文件路径或 HTTP(S) URL。配置自定义模板后不需要再设置 template;embedded_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_template和key_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,为什么没有立即停止?
只有大于 1 的 duration 才会建立截止时间,并且计时从 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 的任务分片逻辑;需要调整生成量时,应结合 frequency 和 objects 评估单个 Task 的批量与执行节奏。