Skip to main content

概述

Spooldir Source Connector 从 Connect Worker 可访问的本地或共享目录发现文件,将文件内容转换为 Kafka 消息并写入指定 Topic。它适合把批量导出的文件、应用日志或落盘交换文件接入 Kafka;每个 Task 扫描输入目录,按文件名筛选文件,读取其中的记录,并在处理成功或失败后按配置移动、删除或保留源文件。 该插件包提供七种 Source Connector 实现。SpoolDirCsvSourceConnector 读取 CSV,并可使用显式 Schema 或从样本文件生成 Schema;SpoolDirJsonSourceConnector 读取连续的 JSON root value,并按显式或生成的 Schema 输出强类型记录;SpoolDirSchemaLessJsonSourceConnector 将连续 JSON root value 输出为无 Schema 的字符串;SpoolDirLineDelimitedSourceConnector 将每一行输出为字符串;SpoolDirBinaryFileSourceConnector 将整个文件作为一条字节数组消息;SpoolDirAvroSourceConnector 读取 Avro 容器文件并使用文件内 Schema;SpoolDirELFSourceConnector 解析 Extended Log Format 文件并动态构造记录 Schema。 七种实现共享文件发现、文件筛选、处理标记和清理策略,但记录结构、时间戳、批次边界和重启后的 Offset 恢复行为并不相同。CSV 是本文快速开始的主路径;为其他格式选择 connector.class 时,只应使用该实现实际支持的配置,不能沿用 CSV 专属配置。

前置条件

  • 若直接解压并部署 Confluent Hub 提供的 2.0.71 官方 ZIP,请在同一插件目录中提供 com.google.guava:guava:31.1-jre,或使用能够将该运行时依赖解析到插件类加载器的安装方式。该 ZIP 未包含 Guava;缺少此依赖时,插件发现或 Task 启动可能失败。补充依赖后需重启 Connect Worker。
  • Connect Worker 必须能以读写权限访问 input.patherror.path;使用 MOVEMOVEBYDATE 清理策略时,还必须能访问并写入 finished.path。这些目录应在 Connector 启动前创建。
  • 使用 CSV 或强类型 JSON 的 Schema 自动生成功能时,Connector 启动前必须在 input.path 顶层放置至少一个匹配 input.file.pattern 的普通文件,并确保用于采样的文件能够生成一致的 Schema。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka、接收数据的 Topic,以及 Worker 可读写的输入、完成和错误目录,并确认文件系统访问权限。具体准备和管理操作请参阅 管理 Connector
将占位符替换为实际资源。启动 Connector 前,在 <input-directory> 顶层放置至少一个以 .csv 结尾的文件;第一行必须是表头,并包含 id 列。Connector 会从匹配文件生成字符串字段 Schema,将 id 写入消息 Key,并在使用默认 MOVE 策略成功处理文件后将其移到 <finished-directory>;无法处理的文件会移到 <error-directory>

配置

输出与文件发现

topic

接收文件记录的 Kafka Topic。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须是 Connector 可写入的有效 Topic 名称。空字符串不会被插件配置校验拒绝,但不能作为可用目标。
  • 必填:是

input.path

待处理文件所在的输入目录。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:目录必须已存在,并且 Connect 进程必须具有读写权限。默认只扫描直接子文件。
  • 必填:是

input.file.pattern

用于选择输入文件名的 Java 正则表达式。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:按完整文件名匹配,不匹配完整路径;不能为空且必须是有效正则表达式。单次匹配超过 100 毫秒时按不匹配处理。
  • 必填:是

file.minimum.age.ms

文件从最后修改到可被发现所需的最短时间,单位毫秒。
  • 类型long
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:必须大于或等于 0。该条件只在发现文件时检查,不会锁定文件,也不会再次确认文件是否仍在写入。

input.path.walk.recursively

是否递归扫描 input.path 下的后代文件。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:设为 true 时仍按每个文件的文件名应用 input.file.pattern

files.sort.attributes

候选文件进入 Task 队列前使用的排序属性。
  • 类型list
  • 默认值NameAsc
  • 重要级别:低
  • 有效值 / 注意事项:每项可为 NameAscNameDescLengthAscLengthDescLastModifiedAscLastModifiedDesc,并按列表顺序组合比较。排序只约束单个 Task 的本地文件队列,不提供跨 Task 全局顺序。

empty.poll.wait.ms

连续空读取后 Task 的等待时间,单位毫秒。
  • 类型long
  • 默认值500
  • 重要级别:低
  • 有效值 / 注意事项:范围为 1Long.MAX_VALUE。它控制空读取后的休眠,不保证固定的文件发现周期。

文件生命周期与错误处理

cleanup.policy

文件成功处理后的清理方式。
  • 类型string
  • 默认值MOVE
  • 重要级别:中
  • 有效值 / 注意事项:可选 NONEDELETEMOVEMOVEBYDATEMOVEMOVEBYDATE 都要求设置可写的 finished.path;处理失败的文件始终尝试移到 error.path

finished.path

成功文件的目标目录。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:使用 MOVEMOVEBYDATE 时必须设置,目录必须在 Task 启动前存在且可写。

error.path

处理失败文件的目标目录。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:目录必须已存在且可写。即使 halt.on.error=false,失败文件也会先尝试移到该目录。
  • 必填:是

halt.on.error

文件处理失败后是否使 Task 失败。
  • 类型boolean
  • 默认值true
  • 重要级别:高
  • 有效值 / 注意事项true 会在错误文件清理后抛出异常;false 会尝试将文件移到 error.path,然后继续选择其他文件。错误处理粒度是整个文件,不是单条记录。

processing.file.extension

处理标记文件使用的后缀。
  • 类型string
  • 默认值.PROCESSING
  • 重要级别:低
  • 有效值 / 注意事项:会追加到输入文件名后;值必须包含点号和非空后缀。存在对应标记时,输入文件不会再次进入候选队列。

cleanup.policy.maintain.relative.path

递归处理后是否保留输入目录中的相对子目录。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只与递归扫描相关。true 保留源端子目录;false 允许清理移动或删除后留下的空子目录。目标目录仍会按实现计算的相对路径创建层级。

file.buffer.size.bytes

读取文件时使用的缓冲区大小,单位字节。
  • 类型int
  • 默认值131072
  • 重要级别:低
  • 有效值 / 注意事项:必须至少为 1。Binary 实现仍会把解压后的整个文件放入单条记录,该配置不是整文件内存上限。

批次、任务分配与时间戳

batch.size

单次处理调用的记录数量上限。
  • 类型int
  • 默认值1000
  • 重要级别:低
  • 有效值 / 注意事项:CSV、强类型 JSON、行文本、Avro、无 Schema JSON 和 ELF 将其作为批次记录数参考;Binary 每个文件只产生一条整文件记录,不按该值拆分内容。不同格式的边界行为存在差异。

task.partitioner

文件分配给 Task 的分区方式。
  • 类型string
  • 默认值ByName
  • 重要级别:中
  • 有效值 / 注意事项:此版本只接受 ByName。多 Task 根据文件 basename 的哈希分配文件,不提供按文件大小的动态均衡;不要设置内部的 task.indextask.count

timestamp.mode

CSV 和强类型 JSON 记录的时间戳来源。
  • 类型string
  • 默认值PROCESS_TIME
  • 重要级别:中
  • 有效值 / 注意事项:可选 FIELDFILE_TIMEPROCESS_TIME。只有 CSV 和 SpoolDirJsonSourceConnector 实际应用该配置;Avro、Binary、行文本、无 Schema JSON 和 ELF 虽注册此键,但输出记录的时间戳为空。

CSV 与强类型 JSON Schema

key.schema

CSV 或强类型 JSON 消息 Key 的 Kafka Connect Schema JSON。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:只适用于 CSV 和 SpoolDirJsonSourceConnector。关闭 Schema 生成时,必须与 value.schema 一起提供;非空值必须能反序列化为 Kafka Connect Schema。

value.schema

CSV 或强类型 JSON 消息 Value 的 Kafka Connect Schema JSON。
  • 类型string
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:只适用于 CSV 和 SpoolDirJsonSourceConnector。关闭 Schema 生成时,必须与 key.schema 一起提供;CSV 列和 JSON 字段按该 Schema 转换为强类型值。

schema.generation.enabled

是否从输入样本生成 CSV 或强类型 JSON 的 Key 和 Value Schema。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只适用于 CSV 和 SpoolDirJsonSourceConnector。启用后,启动时必须在输入目录顶层存在匹配文件;最多采样五个文件,且样本必须生成一致的 Schema。生成字段为可选字符串,不推断数值、布尔或时间类型。

schema.generation.key.fields

Schema 生成时纳入消息 Key 的字段名。
  • 类型list
  • 默认值:空列表
  • 重要级别:中
  • 有效值 / 注意事项:只在 CSV 或强类型 JSON 启用 Schema 生成时使用;指定字段必须存在于生成的 Value Schema 中。

schema.generation.key.name

自动生成的 Key Schema 名称。
  • 类型string
  • 默认值com.github.jcustenborder.kafka.connect.model.Key
  • 重要级别:中
  • 有效值 / 注意事项:只适用于 CSV 和强类型 JSON;启用 Schema 生成时必须解析为非空值。

schema.generation.value.name

自动生成的 Value Schema 名称。
  • 类型string
  • 默认值com.github.jcustenborder.kafka.connect.model.Value
  • 重要级别:中
  • 有效值 / 注意事项:只适用于 CSV 和强类型 JSON;启用 Schema 生成时必须解析为非空值。

parser.timestamp.timezone

解析 CSV 或强类型 JSON 日期时间文本时使用的时区。
  • 类型string
  • 默认值UTC
  • 重要级别:低
  • 有效值 / 注意事项:未知时区 ID 可能被 Java 静默回退为 GMT;应显式使用有效的 Java 时区 ID。

parser.timestamp.date.formats

解析 CSV 或强类型 JSON 时间戳时依次尝试的日期时间格式。
  • 类型list
  • 默认值yyyy-MM-dd'T'HH:mm:ss,yyyy-MM-dd' 'HH:mm:ss
  • 重要级别:低
  • 有效值 / 注意事项:每项必须是有效的 SimpleDateFormat 模式;按列表顺序尝试,建议把更精确的格式放在前面。

timestamp.field

使用字段时间戳时的字段名。
  • 类型string
  • 默认值:空字符串
  • 重要级别:中
  • 有效值 / 注意事项:只适用于 CSV 和强类型 JSON。timestamp.mode=FIELD 时,必须指向 Value Schema 中存在、非可选且逻辑名称为 org.apache.kafka.connect.data.Timestamp 的字段。

CSV 解析与字段映射

csv.skip.lines

解析 CSV 前跳过的起始行数。
  • 类型int
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV。该版本未校验非负范围,应使用大于或等于 0 的值。

csv.separator.char

CSV 分隔符的 Unicode 数值。
  • 类型int
  • 默认值44
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV;44 表示逗号,9 表示制表符。设为 0 会选择 RFC 4180 解析器分支。

csv.quote.char

CSV 引号字符的 Unicode 数值。
  • 类型int
  • 默认值34
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV;34 表示双引号。

csv.escape.char

CSV 转义字符的 Unicode 数值。
  • 类型int
  • 默认值92
  • 重要级别:低
  • 有效值 / 注意事项:只适用于默认 CSV 解析器;92 表示反斜杠。RFC 4180 解析器会忽略此配置。

csv.strict.quotes

默认 CSV 解析器是否只接受引号内的字符。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只适用于默认 CSV 解析器;RFC 4180 解析器会忽略此配置。

csv.ignore.leading.whitespace

默认 CSV 解析器是否忽略引号前的前导空白。
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:只适用于默认 CSV 解析器;RFC 4180 解析器会忽略此配置。

csv.ignore.quotations

默认 CSV 解析器是否忽略引号的特殊含义。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只适用于默认 CSV 解析器;RFC 4180 解析器会忽略此配置。

csv.keep.carriage.return

CSV Reader 是否保留行中的回车字符。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV,并同时作用于默认和 RFC 4180 解析器分支。

csv.verify.reader

CSV Reader 是否验证底层 Reader 状态。
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV,并同时作用于默认和 RFC 4180 解析器分支。

csv.null.field.indicator

将空分隔字段或空引号字段识别为 null 的方式。
  • 类型string
  • 默认值NEITHER
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV;可选 EMPTY_SEPARATORSEMPTY_QUOTESBOTHNEITHER

csv.first.row.as.header

是否将第一条解析记录作为字段名表头。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:只适用于 CSV。设为 true 时表头名称必须与 Value Schema 字段精确匹配;设为 false 时按 value.schema 字段顺序映射列。

csv.file.charset

CSV 文件的字符集。
  • 类型string
  • 默认值Charset.defaultCharset().name()
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV,值必须是 Java 支持的字符集名称。运行时读取使用该值,但 Schema 自动生成仍使用 JVM 默认字符集;非默认编码应优先使用显式 Schema,并在 Worker 上明确默认字符集。

csv.case.sensitive.field.names

声明 CSV 字段名是否区分大小写。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV。此配置在 2.0.71 中虽被注册和读取,但不会改变表头查找或 Schema 生成;实际表头匹配仍区分大小写。

csv.rfc.4180.parser.enabled

是否使用 RFC 4180 CSV 解析器。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:只适用于 CSV。启用后,csv.escape.charcsv.strict.quotescsv.ignore.leading.whitespacecsv.ignore.quotations 不参与该解析器构建。

行文本与无 Schema JSON 字符集

file.charset

行文本或无 Schema JSON 文件声明的字符集。
  • 类型string
  • 默认值Charset.defaultCharset().name()
  • 重要级别:低
  • 有效值 / 注意事项:只注册于 SpoolDirLineDelimitedSourceConnectorSpoolDirSchemaLessJsonSourceConnector。行文本实现按该字符集读取;无 Schema JSON 实现在 2.0.71 中未将该值传给 JSON 解析器。CSV 应使用 csv.file.charset,其他四类不注册字符集配置。

最佳实践

避免读取仍在写入的文件

适用业务场景:上游应用先写出较大的 CSV 文件,再交给 Connector 处理。为了降低 Connector 在文件尚未稳定时开始读取的风险,应让文件达到最短静置时间,并让上游使用不匹配的临时文件名写入,完成后再原子重命名为 .csv 在快速开始配置基础上增加:
关键说明:示例要求文件最后修改时间至少保持 60 秒不变,但该参数只在发现阶段检查,不能替代原子交付。上游可先写入例如 orders.csv.tmp,完成并关闭文件后再重命名为 orders.csv;临时名称不应匹配快速开始中的 input.file.pattern。应按文件生成耗时和可接受接入延迟调整静置时间。

隔离坏文件并继续处理后续文件

适用业务场景:目录会持续接收多个独立文件,单个格式错误或 Schema 不匹配的文件不应让整个 Task 长时间停止。此时可将错误文件移入隔离目录,并让 Task 继续选择其他文件。 在快速开始配置基础上增加:
关键说明:error.path 仍必须存在且可写。错误处理粒度是整个文件;文件中已经发送的记录不会因后续某条记录失败而回滚。应监控错误目录和 Task 日志,修正文件后使用新的唯一文件名重新投放,避免复用历史 basename 和 Offset。

监控

监控内容

关注 Kafka Connect 健康状态、Connector 和 Task 状态、吞吐、延迟、Offset 提交、错误、重试和 Worker JVM 信号;仅在启用了相应错误处理时关注 DLQ 活动。

导入 Grafana 大盘

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

限制条件

  • 七种实现不共享完整的配置和记录语义。Schema、CSV、字符集和时间戳配置必须按所选 connector.class 使用;不能把 CSV 的行为推广到其他格式。
  • Source Partition 只使用文件 basename,不包含目录、Topic、Connector 类或内容摘要。递归目录中的同名文件以及以后重新投放的同名文件会复用历史 Offset,可能导致跳过或重发记录。
  • CSV、强类型 JSON 和 ELF 会使用已保存 Offset 继续读取;行文本、无 Schema JSON 和 Binary 会忽略恢复 Offset,并从文件开头重新发送。Avro 在恢复边界存在重复最后已提交记录的风险,因此不能对七种实现统一承诺可靠断点恢复。
  • 文件清理与 Kafka 写入、Offset 提交之间没有原子协议。该 Connector 不能据此声明 exactly-once、at-least-once、无重复或无丢失;MOVEMOVEBYDATEDELETENONE 在故障窗口中的结果不同。
  • Worker 被强制终止时可能遗留 processing.file.extension 标记,插件没有自动过期或启动清理机制;遗留标记会持续阻止对应文件再次被发现。
  • 递归扫描计算的相对路径包含文件名,移动结果可能形成 <子目录>/<文件名>/<文件名>file.relative.path Header 也可能包含文件名,而不只是父目录。
  • Binary 实现会把解压后的整个文件读入内存并生成一条 Kafka 消息,受 Worker 内存、Producer 请求大小和 Kafka 消息大小限制。
  • timestamp.mode 只对 CSV 和强类型 JSON 生效;其余五种实现创建的记录时间戳为空。
  • ELF 实现在批次达到 batch.size 时会预读下一条日志项但不将其加入结果,存在批次边界记录缺失风险。
  • Schema 自动生成只查看 input.path 顶层的匹配文件,不使用递归扫描;样本只生成可选字符串字段,不能代表后续文件的完整类型约束。

常见问题

Connector 启动时提示输入、完成或错误目录不可用

目录不存在、不是目录或 Connect 进程没有读写权限都会导致 Task 启动失败。确认 input.patherror.path 已创建且可写;使用 MOVEMOVEBYDATE 时还要确认 finished.path 已创建且可写。若目录来自共享卷,还应在实际运行 Worker 的节点或容器内检查挂载路径和权限。

启用 Schema 自动生成后 Connector 无法启动

输入目录顶层没有匹配文件、样本文件无法解析、多个样本生成的 Schema 不一致,或 schema.generation.key.fields 指定了不存在的字段,都可能导致启动失败。启动前在 input.path 顶层放置至少一个符合 input.file.pattern 的普通文件,确保最多前五个样本结构一致,并核对 Key 字段名称。嵌套目录中的样本不会被生成器发现;需要递归读取时,可先用顶层样本启动,或关闭自动生成并显式配置 key.schemavalue.schema

文件已放入输入目录但一直没有被处理

常见原因包括文件名未被 input.file.pattern 完整匹配、文件年龄小于 file.minimum.age.ms、文件位于子目录但未启用递归扫描,或同名的处理标记仍然存在。检查实际文件名、最后修改时间和 input.path.walk.recursively,并查找 <文件名><processing.file.extension>。只有在确认没有 Task 正在处理该文件后,才能人工移除遗留标记。

重启后为什么出现重复记录或跳过记录

不同格式的 Offset 恢复能力不同,而且 Source Partition 只由 basename 标识。行文本、无 Schema JSON 和 Binary 会从头读取重新发现的文件;同名新文件或不同子目录中的同名文件也会复用旧 Offset。为每次投放使用唯一 basename,保留完成目录中的原文件用于核对,并避免在同一 Connector Offset 命名空间中用新内容替换同名文件。对重复或缺失不可接受的链路,应在下游使用稳定业务键进行幂等处理和完整性校验。