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

# Elasticsearch Source Connector

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

## 概述

Elasticsearch Source Connector 将 Elasticsearch 索引中搜索可见的文档读取到 Kafka，适用于索引数据分发、分析数据接入和文档状态同步。每个索引对应一个 Topic，Topic 名称由指定前缀与索引名直接拼接。

Connector 通过 Search API 轮询文档。首次没有保存的读取位置时读取当前文档，随后按业务游标升序读取超过水位的数据。这不是变更日志订阅或 CDC，也不是一致性快照：更新需要推进游标才可能再次读取，删除不会生成事件。消息值是带 Schema 的文档结构，附加文档标识和索引名称；消息键由索引名及游标值组成，而不是稳定的文档 ID。

## 前置条件

* Elasticsearch 索引必须保留可读取的 `_source`，文档中提供非空、顶层、可排序且可执行范围查询的游标字段；单字段游标应唯一，并让后续新增及需要同步的更新严格超过已读取水位。
* Elasticsearch 访问账号需能够查询所选索引，并访问索引发现接口 `/_cat/indices`；即使使用固定索引列表，Connector 仍会访问该接口。
* 使用自定义 TLS 材料时，PKCS12 文件必须在所有可能运行 Connector 和 Task 的 Worker 上可读；配置密码前须控制配置初始化日志输出及日志访问权限，避免密码泄露。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka 和包含游标字段的 Elasticsearch 索引，确认网络连通和访问权限；创建与管理操作参见[管理 Connector](../manage-connectors)。以下示例用于无需认证的受控环境，索引文档使用顶层数值字段 `change_seq` 作为唯一递增游标。

```properties theme={null}
connector.class=com.github.dariobalinzo.ElasticSourceConnector
es.host=<elasticsearch-host>
es.port=9200
index.names=<index-name>
incrementing.field.name=change_seq
topic.prefix=es-
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
```

替换主机和索引名称，确保索引的 `change_seq` 映射支持排序与范围查询，且文档实际含有该字段。主机不包含协议、端口或路径。`es-` 与索引名直接组成输出 Topic。没有历史读取位置时会读取当前可见文档；后续新增或更新必须推进游标并在 Elasticsearch 中搜索可见。示例显式指定 JsonConverter，未设置 `value.converter.schemas.enable`，因此采用该 Converter 的默认值 `true`，输出包含 `schema` 和 `payload`，不继承 Worker 的该子配置。默认字段名称转换会把 `es-id`、`es-index` 变为 `esid`、`esindex`。

## 配置

### Elasticsearch 连接

#### `es.host`

指定 Elasticsearch 主机。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：不带协议、端口或路径；多个主机以分号分隔，共享端口和协议。分隔项不会自动去除空格。

#### `es.port`

指定所有 Elasticsearch 主机使用的端口。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：启动时解析为整数，应填写 `1` 至 `65535` 的合法 TCP 端口；没有显式范围校验。

#### `es.scheme`

指定连接协议。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`http`
* **重要级别**：中
* **有效值 / 注意事项**：使用 `http` 或 `https`；TLS 文件配置不会自动切换为 HTTPS。生产环境使用用户名密码认证时应使用 HTTPS。

#### `es.user`

指定 HTTP Basic 认证用户名。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：高
* **有效值 / 注意事项**：非空时启用用户名密码认证，应配套提供密码；空值关闭此认证。将连接目标限制为可信主机。

#### `es.password`

指定 HTTP Basic 认证密码。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：高
* **有效值 / 注意事项**：仅在用户名非空时使用。此项不是自动隐藏的密码类型，初始化 INFO 日志可能输出解析后的密码。投入生产前必须限制 `com.github.dariobalinzo.ElasticSourceConnectorConfig` 和 `com.github.dariobalinzo.task.ElasticSourceTaskConfig` 的初始化配置日志，并检查日志访问权限；Config Provider 变量替换不能消除此风险。不要在配置示例、排障材料或共享日志中填入真实凭据。

#### `es.tls.truststore.location`

指定自定义 CA 信任库路径。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：仅支持 PKCS12，文件须在相关 Worker 上可读；配合 `https` 使用。不设置此项不等于关闭客户端默认 TLS 校验。

#### `es.tls.truststore.password`

指定信任库密码。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：加载信任库时必须与文件密码匹配，不能为 `null`。同样存在初始化配置日志泄露风险，须先实施密码日志防护。

#### `es.tls.keystore.location`

指定客户端证书及私钥所在的密钥库路径。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：仅支持 PKCS12；只有同时设置信任库路径才会加载客户端密钥库。文件须在相关 Worker 上可读。

#### `es.tls.keystore.password`

指定客户端密钥库密码。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：同一密码用于打开密钥库及读取私钥，不支持独立私钥密码。存在初始化配置日志泄露风险，须先实施密码日志防护。

#### `connection.attempts`

指定搜索发生 I/O 错误时的总尝试次数。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`3`
* **重要级别**：低
* **有效值 / 注意事项**：解析为整数，必须大于 `0`；仅用于搜索的 I/O 异常重试，不覆盖所有异常或索引发现请求。

#### `connection.backoff.ms`

指定搜索 I/O 重试间隔，单位毫秒。

* **类别**：Elasticsearch 连接
* **类型**：`string`
* **默认值**：`10000`
* **重要级别**：低
* **有效值 / 注意事项**：解析为长整数，必须非负，通常使用正值；重试等待会延长数据读取延迟。

### 索引选择与读取位置

#### `index.names`

指定固定索引列表。

* **类别**：索引选择与读取位置
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：逗号分隔的明确索引名，不使用通配符，不自动去除空格。只要提交配置包含此键，就优先使用固定列表；不用固定列表时应移除该项，不能用空字符串或显式 `null` 代替省略。

#### `index.prefix`

按名称前缀选择索引。

* **类别**：索引选择与读取位置
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：单个字面量前缀，不是列表、正则或 glob；空字符串匹配所有索引名。仅在省略 `index.names` 时用于选择，监控线程仍会使用此前缀。

#### `incrementing.field.name`

指定主游标字段。

* **类别**：索引选择与读取位置
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：默认值可以通过配置解析，但不能形成有效查询，实际使用必须指定字段。优先采用 `_source` 中非空顶层数值字段；单游标必须唯一且后续数据严格超过水位。不要把嵌套过滤路径、`_id` 或注入的 `es-id` 直接当作可用游标，也不要直接套用单游标 `.keyword` 路径。

#### `incrementing.secondary.field.name`

指定主游标相同时的次级排序字段。

* **类别**：索引选择与读取位置
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：非 `null` 即启用，空字符串不会关闭。字段须可排序、可范围查询并在每条文档中可取值；主字段应适合精确匹配，组合值应唯一且后续组合严格超过水位。双游标不能补读已越过水位的迟到数据。

#### `mode`

保留的模式选择配置，不改变当前读取方式。

* **类别**：索引选择与读取位置
* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：接受空字符串、`bulk`、`timestamp`、`incrementing`、`timestamp+incrementing`，但没有对应执行分支；不建议设置，不能据此启用独立时间戳或批量模式。此项没有正式弃用标记。

### 输出文档

#### `topic.prefix`

指定输出 Topic 名称前缀。

* **类别**：输出文档
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：允许空字符串；与索引名直接拼接，不自动插入分隔符。修改前缀不会重置读取位置或自动迁移历史数据。

#### `fieldname_converter`

指定文档字段及 Schema 名称的转换方式。

* **类别**：输出文档
* **类型**：`string`
* **默认值**：`avro`
* **重要级别**：中
* **有效值 / 注意事项**：精确小写 `nop` 保留名称，其他值均使用 Avro 名称转换。转换会为非字母开头名称加 `avro` 前缀并移除非 ASCII 字母数字字符，可能造成名称碰撞；不等于启用 Avro 序列化。

#### `filters.whitelist`

仅保留指定文档字段。

* **类别**：输出文档
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：分号分隔路径，支持点分对象路径，不自动去除空格；非 `null` 即启用。需要元数据时显式保留 `es-id` 和 `es-index`；过滤只影响值，不改变消息键及读取位置。

#### `filters.blacklist`

按黑名单规则移除文档字段。

* **类别**：输出文档
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：分号分隔路径，非 `null` 即启用；在白名单之后执行。对象和列表只有存在以该路径开头的黑名单项才会保留，因此未匹配的复合字段也可能被移除，并非通常的仅排除指定叶字段规则。

#### `filters.json_cast`

将指定字段值序列化为 JSON 字符串。

* **类别**：输出文档
* **类型**：`string`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：分号分隔路径，非 `null` 即启用；在白名单和黑名单之后执行。不是解析已有 JSON 字符串，转换字符串值会保留其 JSON 引号。

### 运行与序列化

#### `poll.interval.ms`

指定一次轮询完全没有数据后的等待时间，单位毫秒。

* **类别**：运行与序列化
* **类型**：`string`
* **默认值**：`5000`
* **重要级别**：高
* **有效值 / 注意事项**：解析为整数，必须非负；持续有数据时不按此值限速，也不改变固定为 5000 毫秒的索引发现周期。

#### `batch.max.rows`

指定每次对每个索引搜索的最大文档数。

* **类别**：运行与序列化
* **类型**：`string`
* **默认值**：`10000`
* **重要级别**：低
* **有效值 / 注意事项**：解析为整数，应为正数且不超过 Elasticsearch 查询窗口限制；不是整个 Task 单次轮询的总记录上限，多索引结果会累加。

#### `connector.class`

指定 Source Connector 实现类。

* **类别**：运行与序列化
* **类型**：`string`
* **默认值**：无，必填
* **重要级别**：高
* **有效值 / 注意事项**：使用 `com.github.dariobalinzo.ElasticSourceConnector`。

#### `tasks.max`

指定任务数上限。

* **类别**：运行与序列化
* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：至少为 `1`；实际任务数最多为索引数与此值中的较小者，单索引不会拆分给多个 Task。

#### `key.converter`

指定消息键的序列化 Converter。

* **类别**：运行与序列化
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未在 Connector 级别指定时使用 Worker 配置；键为字符串，快速开始显式使用 StringConverter，该类不是默认值。

#### `value.converter`

指定带 Schema 的 Struct 消息值的序列化 Converter。

* **类别**：运行与序列化
* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：未在 Connector 级别指定时使用 Worker 配置；快速开始显式使用 JsonConverter，该类不是默认值。显式指定此项时，Converter 只使用 Connector 中的 `value.converter.*` 子配置，未设置的子项采用该 Converter 自身的默认值，不合并 Worker 的子配置。JsonConverter 的 `schemas.enable` 默认值为 `true`，快速开始省略此默认设置；字段名称转换不会自动启用 Schema Registry。

## 最佳实践

### 运行期间调整等待时间与批量大小

**适用业务场景**：已经完成接入，希望缩短空闲后新数据的等待时间，同时控制单次读取对 Elasticsearch 和 Worker 内存的压力。

**配置示例**：在快速开始配置上添加以下两项，其余配置保持不变。

```properties theme={null}
poll.interval.ms=1000
batch.max.rows=1000
```

**关键说明**：示例数值是调参起点，并非通用最优值。缩短空闲等待可能增加查询频率；减小页面可降低每个索引单次返回的数据量，但读取存量需要更多搜索请求。应结合查询耗时、Worker 内存和业务延迟要求调整。持续有数据时，等待时间不会限制吞吐；同一轮询中多个索引的结果量会累加。减小页面不能解决重复游标问题，必须先保证游标唯一和推进顺序。

### 扩展到持续创建的滚动索引

**适用业务场景**：数据按日期或周期写入新建索引，维护固定索引列表越来越困难。希望自动发现同一业务前缀下的新索引，并在多个索引间分配读取任务。

**配置示例**：以快速开始为基础，移除 `index.names`，添加以下配置；不要同时保留固定列表。

```properties theme={null}
index.prefix=<business-index-prefix>
tasks.max=2
```

**关键说明**：替换业务索引的字面量前缀，不使用星号。启动前至少准备一个匹配索引，所有匹配索引都应满足同样的游标约束。监控发现列表变化后会重新配置 Task；至少有两个索引时才可能使用两个任务，单索引不能借此扩容。新索引没有保存位置时读取当前可见文档；保持索引名、Connector 身份和游标定义稳定，不把同名索引重建视为自动重新初始化。前缀应足够具体，避免接入无关索引。

## 监控

### 监控内容

关注 Kafka Connect 集群健康、Connector 与 Task 状态、输入输出吞吐、端到端延迟、Offset 提交、错误与重试，以及 Worker JVM 内存、GC 和线程信号；Task 处于 RUNNING 不等于数据持续流动，应结合源端变化和目标消息观察。只有部署启用了相应错误处理时才关注 DLQ 活动。

### 导入 Grafana 大盘

下载共享的 [AutoMQ Connect Cluster Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，确认 Grafana 数据源能够查询 Kafka Connect 和 Worker JVM 指标，并与大盘所需的集群、Connector、Task 等标签匹配，再通过 Grafana 的导入功能加载 JSON 并选择对应数据源。

## 限制条件

* 基于独立搜索读取当前可见状态，不提供 PIT、scroll、一致性快照、变更历史或中间更新版本的完整捕获。
* 单游标按严格大于水位续读；重复游标跨页可能遗漏，迟到且不超过水位的数据不会补读，没有回看窗口或内置去重。
* 双游标仍要求唯一组合且后续组合超过水位，不能补读已越过的主字段组中较小次值的数据。
* 不生成删除事件、tombstone 或 before/after 信封；更新只有在游标推进且搜索可见时才可能再次读取。
* 不提供自定义查询或行筛选配置；字段过滤仅处理输出值。
* 同名索引重建不会清除旧读取位置；改变游标定义或 Topic 前缀也不会自动重置位置。
* 单索引不能拆分到多个 Task，不保证跨索引、跨 Task 或下游多分区的全局顺序。
* 默认消息键随游标变化，不是稳定文档 upsert 主键；附加的 `es-id` 和 `es-index` 会覆盖源文档中同名字段。
* Schema 按每条文档推导，null 和空数组不加入字段；嵌套数组、首元素为 null 的数组不受该转换路径支持，同字段类型冲突不会自动协调。
* 不提供无条件无丢失或 exactly-once 保障；恢复可能重读数据，下游需要考虑重复处理。
* 搜索重试仅覆盖 I/O 异常；部分查询或转换异常在 Task 内记录后返回，不会自动进入框架重试或 DLQ。
* TLS 自定义材料仅支持 PKCS12，客户端密钥库必须配合信任库加载；不支持独立私钥密码、API key 或 OAuth 配置。
* 上游声明适用 Elasticsearch 6.x 和 7.x，不能据此推定 Elasticsearch 8.x 或 9.x 兼容。

## 常见问题

### Task 正常运行，但新增或更新没有消息？

先确认文档已经搜索可见，再检查游标是否缺失、重复，或小于等于已读取水位。不改变游标的更新不会重新读取，双游标也不是迟到补偿机制。检查 Task 错误日志与实际输出，不只依赖运行状态；修正源端游标生成与写入顺序，对已遗漏的数据设计单独回填，不要仅修改 Topic 前缀期待重读。

### 为什么新索引没有自动接入？

检查是否仍存在 `index.names`，固定列表优先且不会自动扩充。需要动态发现时移除该键，确认 `index.prefix` 是匹配实际索引名的字面量前缀，并检查 `/_cat/indices` 权限和错误日志；配置空字符串或星号不能代替正确的前缀选择。

### 重启后如何继续读取，为什么会出现重复？

保留 Connector 身份、索引名、游标定义以及 Kafka Connect 持久化 Offset，才能沿用原读取位置。未持久化的进度在故障后可能重读；没有保存位置时读取当前可见文档，而不是恢复历史删除或中间版本。下游应依据业务标识处理重复；不要将包含游标的默认键视为稳定文档标识。

### 为什么字段名称或 Schema 与源文档不同？

默认名称转换会移除特殊字符，例如附加字段变为 `esid` 和 `esindex`。同时 Schema 来自每条文档的实际值，而不是 Elasticsearch mapping，null、空数组及字段类型变化会影响结果。检查字段过滤顺序和复合字段规则，保持业务字段类型一致，并检查下游是否能接受逐记录 Schema 差异。

### 配置密码后，如何避免日志泄露？

三个密码配置均为字符串类型，初始化配置日志可能打印解析值。上线前验证并限制 Connector 和 Task 配置类的 INFO 配置输出，同时限制日志读取和留存权限；单纯改用 Config Provider 不足以保护解析后的值。发现泄露时立即撤销或轮换凭据，按安全流程处理受影响日志；提交排障材料前必须脱敏，不要再次复制泄露值。
