> ## 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.

# Spooldir Source Connector

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

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector
topic=<topic-name>
input.path=<input-directory>
finished.path=<finished-directory>
error.path=<error-directory>
input.file.pattern=^.*\.csv$
schema.generation.enabled=true
schema.generation.key.fields=id
csv.first.row.as.header=true
```

将占位符替换为实际资源。启动 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`。

在快速开始配置基础上增加：

```properties theme={null}
file.minimum.age.ms=60000
```

关键说明：示例要求文件最后修改时间至少保持 60 秒不变，但该参数只在发现阶段检查，不能替代原子交付。上游可先写入例如 `orders.csv.tmp`，完成并关闭文件后再重命名为 `orders.csv`；临时名称不应匹配快速开始中的 `input.file.pattern`。应按文件生成耗时和可接受接入延迟调整静置时间。

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

适用业务场景：目录会持续接收多个独立文件，单个格式错误或 Schema 不匹配的文件不应让整个 Task 长时间停止。此时可将错误文件移入隔离目录，并让 Task 继续选择其他文件。

在快速开始配置基础上增加：

```properties theme={null}
halt.on.error=false
```

关键说明：`error.path` 仍必须存在且可写。错误处理粒度是整个文件；文件中已经发送的记录不会因后续某条记录失败而回滚。应监控错误目录和 Task 日志，修正文件后使用新的唯一文件名重新投放，避免复用历史 basename 和 Offset。

## 监控

### 监控内容

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

### 导入 Grafana 大盘

确认 Connect 指标已接入 Grafana 数据源，且采集标签满足大盘筛选条件；下载 [Kafka Connect Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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.path` Header 也可能包含文件名，而不只是父目录。
* 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 正在处理该文件后，才能人工移除遗留标记。

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

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