Skip to main content

概述

Shell Source Connector 在运行 Task 的 Kafka Connect Worker 上执行 Shell 命令,将标准输出的每一行转换为一条 Kafka 记录,写入指定 Topic。它连接外部数据源与 Kafka,适合接入已有命令行工具、读取有限文本文件的命令或采集脚本的输出。 每次命令成功结束后,Connector 将该次输出作为一个批次交给 Kafka Connect;随后继续执行下一次采集。记录值是去掉行终止符的字符串,记录键为空,时间戳是记录构造时的时间。采集范围及是否只输出新增数据由命令决定,Connector 不自动识别文件增量或解析 JSON、CSV 字段。

前置条件

  • 实际运行 Task 的每个候选 Worker 操作系统或容器中都必须有可执行的 sh
  • 命令调用的工具、脚本和源资源必须在每个候选 Worker 上可用,并允许 Worker 进程访问;容器中的路径须指向容器内资源。
  • 命令必须能在有限时间内结束,且单次文本输出量应受控,因为整次输出在命令结束后才交给 Kafka Connect。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和目标 Topic,确认网络连通与访问权限,并确认 Worker 可以执行 sh。具体准备和管理操作参见 AutoMQ 的管理 Connector。以下命令使用 Shell 的 printf 输出两行固定文本,无需准备外部文件。
<target-topic> 替换为目标 Topic 名称,并在 Connector 配置中应用上述内容。每次成功采集生成值为 alphabeta 的两条文本记录,默认使用一个 Task;该命令会被重复执行,因此这些文本会持续重复出现,并非只发送一次。代码中的双反斜杠适用于 properties 文件读取,解析后的命令格式字符串为 %s\n

配置

Connector 与执行并发

connector.class

选择 Shell Source Connector 实现。
  • 类别:Kafka Connect 框架配置
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 必填:是
  • 有效值 / 注意事项:使用 uk.co.threefi.connect.shell.ShellSourceConnector

tasks.max

Connector 允许的最大 Task 数。
  • 类别:Kafka Connect 框架配置
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:取值为 12147483647。每个 Task 获得相同命令配置,不会自动划分源数据范围。单源单命令通常保留默认值,增加 Task 可能造成重复执行或资源竞争,而不是实现数据分片。

tasks.max.enforce

是否强制执行 tasks.max 限制。
  • 类别:Kafka Connect 框架配置
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 已弃用:是
  • 替代项:无
  • 有效值 / 注意事项:可取 truefalse。保留默认值,并通过 tasks.max 控制 Task 数;不要依赖关闭此校验来扩展并行度。

命令与采集节奏

shell.command

每次采集执行的 Shell 命令字符串。
  • 类别:Connector 配置
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:应显式填写能成功结束的有效命令。命令由 Worker 上的 sh -c 执行,沿用 Worker 进程环境和工作目录;文件或脚本宜使用绝对路径。仅采集 stdout,按 Worker JVM 的默认字符集解码。命令文本可能出现在日志中,不要写入密码、Token 或其他凭证。没有基于消息键、Topic 或值的模板注入。

block.ms

每次命令处理结束后、返回采集批次前的等待时长,单位毫秒。
  • 类别:Connector 配置
  • 类型int
  • 默认值100
  • 重要级别:低
  • 有效值 / 注意事项:使用 02147483647 的整数。配置解析没有非负范围校验,但实际等待要求非负值。成功、空输出或捕获命令执行 I/O 错误后都会等待;首次采集的输出也在等待后才交给 Connect。它不是命令超时,也不是固定周期或 cron 调度;相邻采集的间隔还包括命令耗时和 Worker 处理时间。

输出目标与文本序列化

topic

所有输出行写入的目标 Kafka Topic。
  • 类别:Connector 配置
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:应显式填写一个有效 Topic 名称,不是名称列表或正则表达式。插件配置校验不检查 Topic 是否有效或是否具备访问权限。

value.converter

指定 Source 记录值写入 Kafka 时使用的 Converter。
  • 类别:Kafka Connect 框架配置
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 的值 Converter。写入纯文本时使用 org.apache.kafka.connect.storage.StringConverter,避免输出格式随 Worker 默认设置变化。Converter 必须是可实例化的 Kafka Connect Converter 实现;选择 Converter 不会让插件自动把 JSON 或 CSV 字符串解析为字段。

最佳实践

持续采集时按源端负载调整等待时间

适用业务场景:首次接入后,需要持续执行每次都能结束的采集命令,但频繁查询或运行脚本会增加源端负载。延长每次采集后的等待时间可以降低访问频率,代价是消息交付延迟增加。 配置示例:在已运行的 Connector 配置中添加以下项,或覆盖已有的 block.ms;保留原来的命令、Topic 和 StringConverter。
关键说明:该示例每次命令完成后等待 3 秒再返回当前批次,不保证每 3 秒执行一次。根据源端可承受的调用频率、命令耗时和业务延迟要求调整等待时长,并观察吞吐和延迟。命令仍应按次结束并控制输出量,不要使用 tail -f 等常驻命令替代周期采集;增加等待不会解决重复输出问题,也不会提供失败后的专用退避策略。

监控

监控内容

关注 Kafka Connect 集群健康、Connector 和 Task 状态、记录吞吐与端到端延迟、Offset 提交成功率及耗时、错误和重试活动,以及 Worker JVM 的内存、GC 和线程信号。Task 状态正常不代表命令每次都成功,应结合源端执行结果与 Kafka 实际输出判断;Offset 提交成功也不代表能够从源位置断点恢复。仅在部署启用了相应错误处理时关注 DLQ 活动,不将命令执行失败视为会自动进入 DLQ 的记录错误。

导入 Grafana 大盘

下载共享的 Kafka Connect Grafana Dashboard,确认监控数据源已采集 Kafka Connect 指标,且集群、Worker、Connector 和 Task 等标签与大盘查询匹配,再通过 Grafana 的导入功能上传 JSON 并选择对应数据源。

限制条件

  • 仅采集 stdout 文本,不自动采集 stderr,也不提供二进制透传。
  • 输出需在命令结束且退出码为零后才返回;持续不结束的命令不能作为逐行实时流使用。
  • 单次输出整体保存在内存中,没有按条数或字节数切批的配置。
  • 不保存或恢复文件位置、行号或源游标,不能依靠已提交的时间戳 Offset 实现断点续读。
  • 不自动去重,也不把源端归档、删除或外部游标推进与 Kafka 写入原子绑定。
  • 不提供消息数据模板注入、内置业务键提取或动态 Topic 路由。
  • 不提供命令执行超时配置或主动终止子进程的插件机制。
  • 不保证跨 Kafka 分区或多个 Task 的全局记录顺序。

常见问题

Task 正常运行,但 Topic 没有新消息怎么办?

确认命令在实际 Worker 容器或主机上、以 Worker 的权限和环境运行时有 stdout 输出,并能结束且退出码为零;检查工具路径、脚本权限和源资源位置。仅写 stderr 的命令不会产生记录,命令输出后以非零状态退出也会丢弃整次输出。查看 Worker 中的命令执行警告,并在不暴露凭证的前提下检查源端错误原因;修复路径、权限或命令后再观察输出。还应核对目标 Topic、Converter 和 block.ms,等待时间过长会延迟当前批次。命令 I/O 错误被捕获后会在后续采集中再次执行,不能只依据 Task 状态判断源端成功。

为什么同一内容会反复出现,重启后也没有从下一行继续?

Connector 每次执行完整命令,不读取历史 Offset 作为源端恢复位置。读取同一个文件或重复查询相同内容都会再次生成记录,重启后也是如此。维护后继续运行前,确认源端数据仍可重放,并让下游按业务标识处理重复;如果业务需要可靠的增量游标和确认后归档,应采用具有相应恢复协议的数据源或 Connector。不要用读完后立即移动、删除文件来代替 Kafka 写入确认,否则可能在尚未送达时失去可重放的数据。

增加 Task 后为什么出现重复采集?

各 Task 执行相同命令,并没有按文件、数据范围或源分区分工。检查 tasks.max,单源单命令恢复为 1,同时检查候选 Worker 是否引用同一资源。不要把增加 Task 数当作源端分片能力;需要扩展时先明确互不重叠的数据范围及资源归属,再设计独立采集任务。

输出 JSON 文本后,为什么消费者仍然收到字符串?

Connector 将每行作为字符串记录,不解析 JSON 字段。使用 StringConverter 时,消费者收到该行文本;若需要结构化数据,可由下游应用解析,或改用能提供结构化记录的数据源 Connector。先检查命令输出的字符集与 Worker JVM 默认字符集是否一致,避免将解码问题误认为 JSON 格式问题。