> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Shell Source Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Shell Source Connector，包括前置条件、配置、监控和故障排查。

## 概述

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](../manage-connectors)。以下命令使用 Shell 的 `printf` 输出两行固定文本，无需准备外部文件。

```properties theme={null}
connector.class=uk.co.threefi.connect.shell.ShellSourceConnector
shell.command=printf '%s\\n' alpha beta
topic=<target-topic>
value.converter=org.apache.kafka.connect.storage.StringConverter
```

将 `<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。

```properties theme={null}
block.ms=3000
```

**关键说明**：该示例每次命令完成后等待 3 秒再返回当前批次，不保证每 3 秒执行一次。根据源端可承受的调用频率、命令耗时和业务延迟要求调整等待时长，并观察吞吐和延迟。命令仍应按次结束并控制输出量，不要使用 `tail -f` 等常驻命令替代周期采集；增加等待不会解决重复输出问题，也不会提供失败后的专用退避策略。

## 监控

### 监控内容

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

### 导入 Grafana 大盘

下载共享的 [Kafka Connect Grafana Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确认监控数据源已采集 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 格式问题。
