概述
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.path和error.path;使用MOVE或MOVEBYDATE清理策略时,还必须能访问并写入finished.path。这些目录应在 Connector 启动前创建。 - 使用 CSV 或强类型 JSON 的 Schema 自动生成功能时,Connector 启动前必须在
input.path顶层放置至少一个匹配input.file.pattern的普通文件,并确保用于采样的文件能够生成一致的 Schema。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka、接收数据的 Topic,以及 Worker 可读写的输入、完成和错误目录,并确认文件系统访问权限。具体准备和管理操作请参阅 管理 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 - 重要级别:低
- 有效值 / 注意事项:每项可为
NameAsc、NameDesc、LengthAsc、LengthDesc、LastModifiedAsc或LastModifiedDesc,并按列表顺序组合比较。排序只约束单个 Task 的本地文件队列,不提供跨 Task 全局顺序。
empty.poll.wait.ms
连续空读取后 Task 的等待时间,单位毫秒。
- 类型:
long - 默认值:
500 - 重要级别:低
- 有效值 / 注意事项:范围为
1到Long.MAX_VALUE。它控制空读取后的休眠,不保证固定的文件发现周期。
文件生命周期与错误处理
cleanup.policy
文件成功处理后的清理方式。
- 类型:
string - 默认值:
MOVE - 重要级别:中
- 有效值 / 注意事项:可选
NONE、DELETE、MOVE或MOVEBYDATE。MOVE和MOVEBYDATE都要求设置可写的finished.path;处理失败的文件始终尝试移到error.path。
finished.path
成功文件的目标目录。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 有效值 / 注意事项:使用
MOVE或MOVEBYDATE时必须设置,目录必须在 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.index或task.count。
timestamp.mode
CSV 和强类型 JSON 记录的时间戳来源。
- 类型:
string - 默认值:
PROCESS_TIME - 重要级别:中
- 有效值 / 注意事项:可选
FIELD、FILE_TIME或PROCESS_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_SEPARATORS、EMPTY_QUOTES、BOTH或NEITHER。
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.char、csv.strict.quotes、csv.ignore.leading.whitespace和csv.ignore.quotations不参与该解析器构建。
行文本与无 Schema JSON 字符集
file.charset
行文本或无 Schema JSON 文件声明的字符集。
- 类型:
string - 默认值:
Charset.defaultCharset().name() - 重要级别:低
- 有效值 / 注意事项:只注册于
SpoolDirLineDelimitedSourceConnector和SpoolDirSchemaLessJsonSourceConnector。行文本实现按该字符集读取;无 Schema JSON 实现在2.0.71中未将该值传给 JSON 解析器。CSV 应使用csv.file.charset,其他四类不注册字符集配置。
最佳实践
避免读取仍在写入的文件
适用业务场景:上游应用先写出较大的 CSV 文件,再交给 Connector 处理。为了降低 Connector 在文件尚未稳定时开始读取的风险,应让文件达到最短静置时间,并让上游使用不匹配的临时文件名写入,完成后再原子重命名为.csv。
在快速开始配置基础上增加:
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、无重复或无丢失;
MOVE、MOVEBYDATE、DELETE和NONE在故障窗口中的结果不同。 - Worker 被强制终止时可能遗留
processing.file.extension标记,插件没有自动过期或启动清理机制;遗留标记会持续阻止对应文件再次被发现。 - 递归扫描计算的相对路径包含文件名,移动结果可能形成
<子目录>/<文件名>/<文件名>,file.relative.pathHeader 也可能包含文件名,而不只是父目录。 - Binary 实现会把解压后的整个文件读入内存并生成一条 Kafka 消息,受 Worker 内存、Producer 请求大小和 Kafka 消息大小限制。
timestamp.mode只对 CSV 和强类型 JSON 生效;其余五种实现创建的记录时间戳为空。- ELF 实现在批次达到
batch.size时会预读下一条日志项但不将其加入结果,存在批次边界记录缺失风险。 - Schema 自动生成只查看
input.path顶层的匹配文件,不使用递归扫描;样本只生成可选字符串字段,不能代表后续文件的完整类型约束。
常见问题
Connector 启动时提示输入、完成或错误目录不可用
目录不存在、不是目录或 Connect 进程没有读写权限都会导致 Task 启动失败。确认input.path 和 error.path 已创建且可写;使用 MOVE 或 MOVEBYDATE 时还要确认 finished.path 已创建且可写。若目录来自共享卷,还应在实际运行 Worker 的节点或容器内检查挂载路径和权限。
启用 Schema 自动生成后 Connector 无法启动
输入目录顶层没有匹配文件、样本文件无法解析、多个样本生成的 Schema 不一致,或schema.generation.key.fields 指定了不存在的字段,都可能导致启动失败。启动前在 input.path 顶层放置至少一个符合 input.file.pattern 的普通文件,确保最多前五个样本结构一致,并核对 Key 字段名称。嵌套目录中的样本不会被生成器发现;需要递归读取时,可先用顶层样本启动,或关闭自动生成并显式配置 key.schema 和 value.schema。
文件已放入输入目录但一直没有被处理
常见原因包括文件名未被input.file.pattern 完整匹配、文件年龄小于 file.minimum.age.ms、文件位于子目录但未启用递归扫描,或同名的处理标记仍然存在。检查实际文件名、最后修改时间和 input.path.walk.recursively,并查找 <文件名><processing.file.extension>。只有在确认没有 Task 正在处理该文件后,才能人工移除遗留标记。