Skip to main content

概述

Shell Sink Connector 从 Kafka Topic 消费记录,并针对每条记录在承载 Task 的 Connect Worker 上启动一个本地 shell 进程。Connector 将 shell.command 中的 ${key}${topic}${value} 替换为当前记录的 Key、Topic 和 Value,再通过 sh -c 执行替换后的完整命令。Key 或 Value 为 null 时,对应占位符替换为空字符串。 单个 Task 按收到的批次顺序逐条执行命令,前一条命令退出后才处理下一条记录;多个 Task 可以在不同 Partition 上并发执行。该 Connector 适合执行依赖受控命令和受控输入的轻量自动化操作。占位符替换不进行 shell 转义,且替换后的完整命令会写入 INFO 日志,因此不适合直接拼接不可信消息或敏感数据。

前置条件

  • 所有可能承载 Task 的 Connect Worker 主机或容器都必须提供 sh,并具备命令所需的可执行文件、固定脚本、文件路径和权限;命令产生的本地文件或其他副作用位于实际运行 Task 的 Worker 上。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、输入 Topic,并确认承载 Task 的 Connect Worker 可以写入 /tmp。具体准备和管理操作请参阅 管理 Connector
<topic-name> 替换为要消费的 Kafka Topic。该固定命令不插入消息 Key 或 Value,每消费一条记录就在当前 Task 所在 Worker 的 /tmp/shell-sink-events.log 中追加一行 processed,可用于确认基本执行链路。替换后的完整命令会写入 Worker 的 INFO 日志;用于实际业务时,应改为预先部署、权限受限且有明确执行时限的固定程序或脚本,不要直接拼接不可信或敏感消息内容。

配置

命令执行

shell.command

针对每条 Kafka 记录执行的 shell 命令模板。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:可使用 ${key}${topic}${value}。Connector 先执行文本替换,再通过 sh -c 运行完整命令;占位符值不会被引用或转义,完整替换结果会写入 INFO 日志。空字符串可以通过配置解析,但不会产生有用的业务副作用。命令在 Task 所在 Worker 的环境、当前目录、文件系统和权限下执行,每条记录启动一个进程,Connector 不提供执行超时。

失败重试

max.retries

shell 命令失败后允许的 Connector 自有重试次数。
  • 类型int
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:最小值为 00 表示首次命令失败后直接使 Task 失败;正整数 N 表示首次尝试之外最多允许 N 次重试。重试会重新执行整个 Worker 批次,而不是只执行失败记录;一次批次成功后,Task 的重试预算恢复为配置值。

retry.backoff.ms

发生可重试命令失败后,传给 Kafka Connect 的下一轮 consumer.poll 超时上限,单位为毫秒。
  • 类型int
  • 默认值3000
  • 重要级别:中
  • 有效值 / 注意事项:最小值为 0,仅在 max.retries 大于 0 且仍有重试预算时使用。在 Kafka Connect 3.9.1 中,该值是下一轮 consumer.poll 允许阻塞的最长时间,还可能被距下一次 Offset 提交的剩余时间进一步缩短;Poll 也可能提前返回,因此不保证实际重投前至少等待该时长。该配置不是 shell 命令执行超时。

Connector 与任务

connector.class

要加载的 Shell Sink Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:使用 uk.co.threefi.connect.shell.ShellSinkConnector
  • 必填:是

tasks.max

允许启动的最大 Sink Task 数量。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:最小值为 1。实际并行度还受输入 Topic Partition 数量和 Worker 分配限制。不同 Task 可以并发执行命令,因此目标脚本和外部系统必须能够处理并发调用;Connector 不提供跨 Partition 或跨 Task 的全局执行顺序。

Kafka 输入

topics

要消费的 Kafka Topic 列表。
  • 类型list
  • 默认值[]
  • 重要级别:高
  • 有效值 / 注意事项:使用逗号分隔多个 Topic。必须在 topicstopics.regex 中恰好配置一个。
  • 必填:条件必填

topics.regex

用于动态匹配输入 Topic 的 Java 正则表达式。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:必须是有效的 Java 正则表达式。必须在 topics.regextopics 中恰好配置一个。
  • 必填:条件必填

记录转换

key.converter

在替换 ${key} 前反序列化 Kafka 消息 Key 的 Converter 类。
  • 类型class
  • 默认值:未设置(null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时继承 Worker 级 Key Converter。Connector 对转换后的 Key 调用 Java toString()null Key 替换为空字符串。Converter 输出不会由 Connector 自动进行 shell 转义。

value.converter

在替换 ${value} 前反序列化 Kafka 消息值的 Converter 类。
  • 类型class
  • 默认值:未设置(null
  • 重要级别:低
  • 有效值 / 注意事项:未设置时继承 Worker 级 Value Converter。Connector 对转换后的值调用 Java toString();墓碑记录的 null 值替换为空字符串。Map、Struct、数组和字节数组等结构化值使用对应 Java 对象的字符串形式,不保证保留原始 JSON、CSV 或二进制表示,也不会自动进行 shell 转义。

最佳实践

临时故障时使用有限重试

适用业务场景: 已完成首次接入,固定脚本偶尔会因短暂资源竞争或目标服务不可用而返回非零退出码。命令对应的外部操作必须能够安全重复执行,并且希望 Task 自动进行有限次数的重试,而不是首次失败就停止。 配置示例:
关键说明:/opt/connect/bin/process-shell-event 作为权限受限的固定脚本部署到所有候选 Worker,并让脚本自行实施明确的执行时限。max.retries=3 表示首次尝试之外最多重试三次;预算耗尽后 Task 失败。retry.backoff.ms=5000 请求将下一轮 Poll 的阻塞超时上限设为 5000 毫秒,但不保证实际重投前至少等待 5000 毫秒。重试粒度是整个 Worker 批次:如果批次前几条记录已经执行成功,后续记录失败会使已成功的批次前缀再次执行。只有脚本和目标系统具备幂等或去重能力时才启用重试;该配置不提供 exactly-once 保证。

积压增长时按 Partition 扩展 Task

适用业务场景: Connector 已稳定运行,输入 Topic 有多个 Partition,单个 Task 的串行命令执行造成持续积压,并且固定脚本及目标系统能够承受多个调用并发执行。可逐步增加 Task 上限,利用 Partition 并行处理记录。 配置示例:
关键说明: tasks.max=4 只是 Task 上限,有效并行度不会超过输入 Partition 数量。每个 Task 内仍按批次顺序逐条执行,但不同 Task 可以在同一或不同 Worker 上并发运行,不保证跨 Partition 的全局顺序。扩展前应确认脚本在每台候选 Worker 上一致可用,目标文件、锁和外部服务支持并发,并根据积压、命令耗时和目标容量逐步调整,而不是把示例值视为固定最优值。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;同时关注长时间无进展的 Task 和外部命令耗时,避免无超时命令长期阻塞消费。仅在启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源,且采集标签满足大盘筛选条件;下载 Kafka Connect Dashboard,在 Grafana 中导入 JSON 并选择对应数据源。

限制条件

  • ${key}${topic}${value} 只是未转义的文本替换,随后由 sh -c 解析;将不可信记录内容直接放入命令可能改变命令语法并造成命令注入。
  • 替换后的完整命令会写入 Worker 的 INFO 日志,命令模板和占位符内容不得包含密码、令牌、凭证或其他敏感数据。
  • Connector 不提供 shell 命令执行超时,也不能在 Task 停止时主动终止正在运行的子进程;不退出的命令会长期阻塞该 Task。
  • Connector 不支持 exactly-once,也不把 shell 副作用与 Kafka Offset 提交放入同一事务。整批重试以及命令成功后、Offset 提交前的故障都可能导致重复执行。
  • 每条记录都会启动独立进程;单个 Task 内命令串行执行,长命令和进程启动开销会直接限制吞吐。
  • ${key}${value} 使用 Converter 输出对象的 Java 字符串形式,Connector 不提供 JSON、CSV、二进制编码或字段路径模板。
  • Kafka Connect 的 errors.tolerance 和 DLQ 机制不能跳过或接收 SinkTask.put 内发生的 shell 命令失败;Connector 只使用 max.retriesretry.backoff.ms 对整个批次进行有限重试,预算耗尽后 Task 失败。

常见问题

命令返回非零退出码后 Task 为什么停止

max.retries 的默认值是 0,因此命令首次失败后会直接产生不可恢复的 Task 错误。检查实际运行 Task 的 Worker 日志、命令退出码、脚本路径、执行权限和依赖程序,并使用与 Worker 相同的用户和环境复现命令。只有故障确实短暂且外部操作可安全重复时,才配置有限的 max.retriesretry.backoff.ms;修正问题后按 Connect 管理流程重启失败的 Task。

为什么目标系统出现重复操作

命令失败后的重试会重新提交整个 Worker 批次,批次中已成功执行的记录可能再次执行;即使整个批次成功,Worker 在提交 Kafka Offset 前发生故障也会重放记录。不要把 Kafka Offset 提交理解为 shell 副作用的事务确认。优先使用幂等命令和目标端去重,并确保去重标识来自经过严格校验、不会改变 shell 语法且不会泄露敏感信息的受控输入。

Task 仍在运行但长时间没有消费进展

某条命令可能长时间运行、等待输入或阻塞输出,而 Connector 没有命令超时和主动终止机制。先在 Task 所在 Worker 上检查对应子进程及脚本日志,再修正固定程序或脚本,使其不依赖交互式输入、限制输出量并实施明确的连接和执行时限。处理遗留进程后重启 Task,并持续观察消费延迟和 Offset 是否恢复推进。

可以用引号安全包裹消息值吗

不能把在模板外层添加单引号或双引号当作通用防护。占位符会先被原样替换,消息中的引号、命令替换、重定向符、换行或其他 shell 元字符仍可能突破预期语法。不要直接插入不可信 Key 或 Value;确需使用记录内容时,应在进入 Connector 前把输入限制为严格的允许字符集合,避免敏感数据,并使用经过安全审查的固定命令边界。