概述
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 配置中应用上述内容。每次成功采集生成值为 alpha 和 beta 的两条文本记录,默认使用一个 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 - 重要级别:高
- 有效值 / 注意事项:取值为
1至2147483647。每个 Task 获得相同命令配置,不会自动划分源数据范围。单源单命令通常保留默认值,增加 Task 可能造成重复执行或资源竞争,而不是实现数据分片。
tasks.max.enforce
是否强制执行 tasks.max 限制。
- 类别:Kafka Connect 框架配置
- 类型:
boolean - 默认值:
true - 重要级别:低
- 已弃用:是
- 替代项:无
- 有效值 / 注意事项:可取
true或false。保留默认值,并通过tasks.max控制 Task 数;不要依赖关闭此校验来扩展并行度。
命令与采集节奏
shell.command
每次采集执行的 Shell 命令字符串。
- 类别:Connector 配置
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:应显式填写能成功结束的有效命令。命令由 Worker 上的
sh -c执行,沿用 Worker 进程环境和工作目录;文件或脚本宜使用绝对路径。仅采集 stdout,按 Worker JVM 的默认字符集解码。命令文本可能出现在日志中,不要写入密码、Token 或其他凭证。没有基于消息键、Topic 或值的模板注入。
block.ms
每次命令处理结束后、返回采集批次前的等待时长,单位毫秒。
- 类别:Connector 配置
- 类型:
int - 默认值:
100 - 重要级别:低
- 有效值 / 注意事项:使用
0至2147483647的整数。配置解析没有非负范围校验,但实际等待要求非负值。成功、空输出或捕获命令执行 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。
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 数当作源端分片能力;需要扩展时先明确互不重叠的数据范围及资源归属,再设计独立采集任务。