Skip to main content

概述

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。以下示例用于无需认证的受控环境,索引文档使用顶层数值字段 change_seq 作为唯一递增游标。
替换主机和索引名称,确保索引的 change_seq 映射支持排序与范围查询,且文档实际含有该字段。主机不包含协议、端口或路径。es- 与索引名直接组成输出 Topic。没有历史读取位置时会读取当前可见文档;后续新增或更新必须推进游标并在 Elasticsearch 中搜索可见。示例显式指定 JsonConverter,未设置 value.converter.schemas.enable,因此采用该 Converter 的默认值 true,输出包含 schemapayload,不继承 Worker 的该子配置。默认字段名称转换会把 es-ides-index 变为 esidesindex

配置

Elasticsearch 连接

es.host

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

es.port

指定所有 Elasticsearch 主机使用的端口。
  • 类别:Elasticsearch 连接
  • 类型string
  • 默认值:无,必填
  • 重要级别:高
  • 有效值 / 注意事项:启动时解析为整数,应填写 165535 的合法 TCP 端口;没有显式范围校验。

es.scheme

指定连接协议。
  • 类别:Elasticsearch 连接
  • 类型string
  • 默认值http
  • 重要级别:中
  • 有效值 / 注意事项:使用 httphttps;TLS 文件配置不会自动切换为 HTTPS。生产环境使用用户名密码认证时应使用 HTTPS。

es.user

指定 HTTP Basic 认证用户名。
  • 类别:Elasticsearch 连接
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:非空时启用用户名密码认证,应配套提供密码;空值关闭此认证。将连接目标限制为可信主机。

es.password

指定 HTTP Basic 认证密码。
  • 类别:Elasticsearch 连接
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:仅在用户名非空时使用。此项不是自动隐藏的密码类型,初始化 INFO 日志可能输出解析后的密码。投入生产前必须限制 com.github.dariobalinzo.ElasticSourceConnectorConfigcom.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
  • 默认值:空字符串
  • 重要级别:高
  • 有效值 / 注意事项:接受空字符串、bulktimestampincrementingtimestamp+incrementing,但没有对应执行分支;不建议设置,不能据此启用独立时间戳或批量模式。此项没有正式弃用标记。

输出文档

topic.prefix

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

fieldname_converter

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

filters.whitelist

仅保留指定文档字段。
  • 类别:输出文档
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:分号分隔路径,支持点分对象路径,不自动去除空格;非 null 即启用。需要元数据时显式保留 es-ides-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 内存的压力。 配置示例:在快速开始配置上添加以下两项,其余配置保持不变。
关键说明:示例数值是调参起点,并非通用最优值。缩短空闲等待可能增加查询频率;减小页面可降低每个索引单次返回的数据量,但读取存量需要更多搜索请求。应结合查询耗时、Worker 内存和业务延迟要求调整。持续有数据时,等待时间不会限制吞吐;同一轮询中多个索引的结果量会累加。减小页面不能解决重复游标问题,必须先保证游标唯一和推进顺序。

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

适用业务场景:数据按日期或周期写入新建索引,维护固定索引列表越来越困难。希望自动发现同一业务前缀下的新索引,并在多个索引间分配读取任务。 配置示例:以快速开始为基础,移除 index.names,添加以下配置;不要同时保留固定列表。
关键说明:替换业务索引的字面量前缀,不使用星号。启动前至少准备一个匹配索引,所有匹配索引都应满足同样的游标约束。监控发现列表变化后会重新配置 Task;至少有两个索引时才可能使用两个任务,单索引不能借此扩容。新索引没有保存位置时读取当前可见文档;保持索引名、Connector 身份和游标定义稳定,不把同名索引重建视为自动重新初始化。前缀应足够具体,避免接入无关索引。

监控

监控内容

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

导入 Grafana 大盘

下载共享的 AutoMQ Connect Cluster Dashboard,确认 Grafana 数据源能够查询 Kafka Connect 和 Worker JVM 指标,并与大盘所需的集群、Connector、Task 等标签匹配,再通过 Grafana 的导入功能加载 JSON 并选择对应数据源。

限制条件

  • 基于独立搜索读取当前可见状态,不提供 PIT、scroll、一致性快照、变更历史或中间更新版本的完整捕获。
  • 单游标按严格大于水位续读;重复游标跨页可能遗漏,迟到且不超过水位的数据不会补读,没有回看窗口或内置去重。
  • 双游标仍要求唯一组合且后续组合超过水位,不能补读已越过的主字段组中较小次值的数据。
  • 不生成删除事件、tombstone 或 before/after 信封;更新只有在游标推进且搜索可见时才可能再次读取。
  • 不提供自定义查询或行筛选配置;字段过滤仅处理输出值。
  • 同名索引重建不会清除旧读取位置;改变游标定义或 Topic 前缀也不会自动重置位置。
  • 单索引不能拆分到多个 Task,不保证跨索引、跨 Task 或下游多分区的全局顺序。
  • 默认消息键随游标变化,不是稳定文档 upsert 主键;附加的 es-ides-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 与源文档不同?

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

配置密码后,如何避免日志泄露?

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