SeaTunnel Easysearch Source Connector 完全指南:从 INFINI Easysearch 批量读取数据
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文围绕 SeaTunnel 的 Easysearch Source 连接器展开,系统讲解如何通过 SeaTunnel 从 INFINI Easysearch 集群批量读取数据。你将掌握连接器的核心能力边界、全部配置参数及其源码级实现原理、Easysearch 与 SeaTunnel 之间的类型映射规则,并拿到可直接复制运行的 HOCON 配置示例,涵盖字段裁剪、DSL 查询过滤、Schema 显式转换以及 HTTPS/TLS 安全连接等实战场景。
一、连接器概述:定位与能力边界
Easysearch Source 连接器用于从 INFINI Easysearch 读取数据,是 SeaTunnel 连接 Easysearch 生态的数据入口。与同一个连接器模块内的 Easysearch Sink 配合,可以实现 Easysearch 集群之间的数据搬迁、索引同步等场景。
从源码结构看,该连接器位于 seatunnel-connectors-v2/connector-easysearch,包含 Source、Sink、Catalog(EasysearchCatalog)三大部分;本文聚焦其中的 Source 部分。
1.1 支持的计算引擎
| 引擎 | 支持情况 |
|---|---|
| Spark | ✔ 支持 |
| Flink | ✔ 支持 |
| SeaTunnel Zeta | ✔ 支持 |
官方声明该连接器支持 INFINI Easysearch 官方发布的所有版本(见 Easysearch.md)。
1.2 关键特性清单
从连接器特性矩阵看,Easysearch Source 的能力边界如下:
- ✅ batch(批式读取)
- ❌ stream(流式读取)
- ❌ exactly-once(精确一次)
- ✅ column projection(列投影)
- ❌ parallelism(并行度)
- ❌ support user-defined split(用户自定义分片)
结合源码可以确认这些特性的含义:EasysearchSource.java 中getBoundedness()返回Boundedness.BOUNDED,即这是一个有界(批式)数据源;同时该类实现了SupportParallelism与SupportColumnProjection接口——需要注意的是,虽然连接器声明支持并行接口,但每个并行子任务都会对同一批索引发起独立的 Scroll 请求,文档特性表中 parallelism 仍标记为未勾选,说明并行能力并未作为成熟特性对外承诺。
1.3 依赖说明
使用该连接器需要引入 Easysearch 官方客户端依赖easysearch-client(文档中给出的 Maven 中央仓库坐标为com.infinilabs:easysearch-client)。在源码中,EasysearchClient.java 正是基于该 RestClient 封装的 HTTP 访问层,统一负责连接管理、认证、TLS 配置与 Scroll 请求。
二、数据读取原理:基于 Scroll 的批量拉取
理解 Easysearch Source 的工作机制,是正确配置参数的前提。从源码可以梳理出完整的读取链路:
1. 分片枚举阶段(由EasysearchSourceSplitEnumerator负责):
- 通过
_cat/indices/{index}?h=index,docsCount&format=json接口(见 EasysearchClient.java)获取索引列表及其文档数; - 过滤掉文档数为 0 的索引,并按文档数升序排序;
- 每个非空索引被包装为一个
EasysearchSourceSplit(分片 ID 为索引名 hashCode); - 分片通过
assignCount % readerCount轮询方式分发给各 Reader(见 EasysearchSourceSplitEnumerator.java)。
2. 数据拉取阶段(由EasysearchSourceReader负责,见 EasysearchSourceReader.java):
- 对每个分片调用
/{index}/_search?scroll={scroll_time}发起首次 Scroll 请求; - 循环调用
/_search/scroll携带 scrollId 继续翻页,直到返回的文档列表为空; - 每次请求完成后,在
finally块中调用clearScroll清理服务端 Scroll 上下文,避免资源泄漏(清理失败仅记录 warn 日志,不影响主流程)。
3. 反序列化阶段(由DefaultSeaTunnelRowDeserializer负责):将每个文档的_source字段按配置的 SeaTunnel RowType 逐字段转换并输出为SeaTunnelRow。
这个设计意味着:Easysearch Source 本质上是「按索引拆分 + Scroll 游标翻页」的批式读取器,index支持通配符时,每个匹配的物理索引都会成为一个独立数据分片。
三、Source 参数详解
以下是 Easysearch Source 的全部参数,源码定义位于 EasysearchSourceOptions.java 与 EasysearchSinkCommonOptions.java:
| 名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| hosts | array | 是 | - | Easysearch HTTP 地址列表 |
| index | string | 是 | - | Easysearch 索引名,支持seatunnel-*通配匹配 |
| username | string | 否 | - | 安全认证用户名 |
| password | string | 否 | - | 安全认证密码 |
| source | array | 否 | - | 需要读取的字段列表;与schema二选一 |
| schema | config | 否 | - | 用于读取与转换字段的 SeaTunnel Schema;与source二选一 |
| query | json | 否 | {"match_all":{}} | Easysearch DSL 查询,用于过滤记录 |
| scroll_time | string | 否 | 1m | Easysearch 保持 Scroll 上下文存活的时间 |
| scroll_size | int | 否 | 100 | 每次 Scroll 请求返回的最大记录数 |
| tls_verify_certificate | boolean | 否 | true | 是否校验 HTTPS 证书 |
| tls_verify_hostname | boolean | 否 | true | 是否校验 HTTPS 主机名 |
| tls_keystore_path | string | 否 | - | PEM 或 JKS 密钥库路径 |
| tls_keystore_password | string | 否 | - | 密钥库密码 |
| tls_truststore_path | string | 否 | - | PEM 或 JKS 信任库路径 |
| tls_truststore_password | string | 否 | - | 信任库密码 |
| common-options | config | 否 | - | Source 插件通用参数 |
3.1 hosts [array](必填)
Easysearch 集群 HTTP 地址,格式为host:port,支持配置多个地址实现故障转移,例如["host1:9200", "host2:9200"]。
在源码中,EasysearchClient.java 会将每个 host 通过HttpHost.create()转换为HttpHost并交给 RestClient 构建器。同时源码设置了两个连接超时值:连接请求超时CONNECTION_REQUEST_TIMEOUT = 10s、Socket 超时SOCKET_TIMEOUT = 5min,数据量较大时可据此评估 Scroll 请求是否会超时。
3.2 index [string](必填)
要读取的 Easysearch 索引名,支持*通配符匹配。如seatunnel-*可匹配所有以seatunnel-开头的索引。通配符场景下,每个匹配的索引会被拆分为独立分片,按文档数升序逐个读取。
3.3 username / password [string](可选)
Easysearch 安全认证的用户名与密码。源码中当username存在时,会通过BasicCredentialsProvider为所有请求域(AuthScope.ANY)设置UsernamePasswordCredentials(见 EasysearchClient.java)。注意:password可以单独缺省(此时按空密码处理),但建议成对配置。
3.4 source [array](可选,与 schema 二选一)
指定要从索引中读取的字段列表。特别地,可以通过指定字段_id获取文档 ID;若后续需要将_id写入另一个 Easysearch 索引,必须为_id指定别名,因为 Easysearch 不允许把_id作为普通字段写入(该限制同样适用于 Easysearch Sink 端)。
当source与schema均未配置时,连接器会调用/{index}/_mappings接口(见 EasysearchClient.java)获取索引 Mapping,自动使用索引中的全部映射字段。
3.5 schema [config](可选,与 source 二选一)
数据的结构定义,包含字段名与字段类型,用于让 SeaTunnel 按显式类型定义转换所选字段。更多细节参见 Schema Feature。
schema与source互斥。当两者都省略时,连接器自动读取 Easysearch 字段映射并使用索引中所有映射字段。从DefaultSeaTunnelRowDeserializer的源码(DefaultSeaTunnelRowDeserializer.java)看,schema支持的能力远不止标量类型:
- Map 类型:如
map<string, tinyint>; - Array 类型:如
array<tinyint>,元素会按元素类型递归转换; - Decimal 类型:如
decimal(2, 1),按BigDecimal解析; - Bytes 类型:按 Base64 解码为字节数组;
- 日期时间类型:
date/timestamp支持毫秒时间戳(Instant.ofEpochMilli+ 系统时区)以及多种字符串格式(如yyyy-MM-dd HH:mm:ss、yyyy-MM-dd HH:mm:ss.SSS等)的自动解析; - 嵌套字段:字段名支持
.分隔的路径(如user.name),通过recursiveGet递归取嵌套值(见 DefaultSeaTunnelRowDeserializer.java)。
提示:当
source指定的字段在索引 Mapping 中不存在时,客户端会回退为默认类型text并输出 warn 日志(见 EasysearchClient.java)。
3.6 query [json](可选,默认{"match_all":{}})
Easysearch 查询 DSL,用于控制读取的数据范围。默认值为{"match_all":{}}(见 EasysearchSourceOptions.java)。在发起首次 Scroll 请求时,该 query 会被放入请求体的query字段(见 EasysearchClient.java),同时请求体固定附带sort: ["_doc"](按内部文档序,Scroll 场景下效率最高)与size: scroll_size。
3.7 scroll_time [string](可选,默认 1m)
Easysearch 为 Scroll 请求保持搜索上下文存活的时间。首次请求通过/_search?scroll={scroll_time}创建上下文,后续每次翻页都会携带该参数续期。该值需大于单次拉取全部数据的耗时,否则上下文过期会导致读取中断;数据量大时建议调大(如2m、5m)。
3.8 scroll_size [int](可选,默认 100)
每次 Scroll 请求返回的最大命中数。增大该值可减少请求往返次数、提升吞吐,但会增大单次响应的内存开销,需结合文档大小权衡。
3.9 TLS 相关参数(可选)
Easysearch Source 完整支持 HTTPS 安全连接,相关选项在 EasysearchClient.java 中有明确实现逻辑:
tls_verify_certificate = true(默认):校验 HTTPS 证书,此时可通过tls_keystore_path/tls_keystore_password/tls_truststore_path/tls_truststore_password配置密钥库与信任库(支持 PEM 与 JKS 格式);配置文件必须对运行 SeaTunnel 的操作系统用户可读;tls_verify_certificate = false:跳过证书校验(源码使用TrustAllStrategy信任所有证书);tls_verify_hostname = false:跳过主机名校验(源码使用NoopHostnameVerifier)。
从源码逻辑看,tls_verify_certificate控制是否构建并加载 SSLContext 及信任策略,tls_verify_hostname独立控制 HostnameVerifier,二者可分别关闭。
3.10 common options
Source 插件通用参数,详见 Source Common Options。
四、数据类型映射
Easysearch Source 内置了 Easysearch 字段类型到 SeaTunnel 类型的映射表,源码实现于 EzsTypeMappingSeaTunnelType.java:
| Easysearch 数据类型 | SeaTunnel 数据类型 |
|---|---|
| STRING / KEYWORD / TEXT | STRING |
| BOOLEAN | BOOLEAN |
| BYTE | BYTE |
| SHORT | SHORT |
| INTEGER | INT |
| LONG | LONG |
| FLOAT / HALF_FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| DATE | LOCAL_DATE_TIME_TYPE |
几点源码层面的补充:
- 映射表不区分大小写场景有限:映射键为小写(
string、keyword、text、integer、half_float等),Mapping 中获取的类型会被直接用于查表; - 源码中还额外支持了
binary类型映射到 STRING(文档表中未单列); - 当 Easysearch 类型不在映射表内时,会抛出
EasysearchConnectorException(错误码EZS_FIELD_TYPE_NOT_SUPPORT),例如nested、object等复合类型若直接参与自动 Mapping 转换可能报错,此时建议用schema显式声明字段类型; - 上述映射主要用于「未配置 source/schema 时自动推导字段类型」的场景(对应测试 EasysearchSourceTest.java 中
testPrepareWithEmptySource的验证逻辑);而一旦配置了schema,则以 schema 声明的 SeaTunnel 类型为准,反序列化器会按该类型执行转换(支持 Decimal、Array、Map、Bytes、日期时间等丰富类型,见上一节)。
五、配置示例
以下示例均来自官方文档,可直接用于 SeaTunnel 的 HOCON 配置文件(.conf)。
5.1 读取指定字段(source 模式)
只读取索引中的指定字段,并通过 DSL 过滤记录:
source { Easysearch { hosts = ["localhost:9200"] index = "seatunnel-*" source = ["_id", "name", "age"] query = {"range": {"age": {"gte": 18, "lte": 60}}} } }该配置的含义:扫描所有seatunnel-*索引,读取每个文档的_id、name、age三个字段,且仅保留age在 18 到 60 之间的记录(DSLrange查询在服务端完成过滤)。
5.2 使用 schema 显式声明字段(schema + query 模式)
当需要 SeaTunnel 按显式类型转换字段时,使用schema声明完整结构。下面的示例还演示了带认证的 HTTPS 连接、关闭证书与主机名校验,以及完整的 Batch 任务(包含 Easysearch Sink 形成搬迁闭环):
env { parallelism = 1 job.mode = "BATCH" } source { Easysearch { hosts = ["https://e2e_easysearch:9200"] username = "admin" password = "admin" tls_verify_certificate = false tls_verify_hostname = false index = "st_index" query = {"range": {"c_int": {"gte": 10, "lte": 20}}} schema = { fields { c_map = "map<string, tinyint>" c_array = "array<tinyint>" c_string = string c_boolean = boolean c_tinyint = tinyint c_smallint = smallint c_int = int c_bigint = bigint c_float = float c_double = double c_decimal = "decimal(2, 1)" c_bytes = bytes c_date = date c_timestamp = timestamp } } } } sink { Easysearch { hosts = ["https://e2e_easysearch:9200"] username = "admin" password = "admin" tls_verify_certificate = false tls_verify_hostname = false index = "st_index2" } }这是一个「Easysearch → SeaTunnel → Easysearch」的完整数据搬迁任务:从st_index按c_int区间过滤读取,schema 声明的字段经类型转换后写入st_index2。由于_id不能作为普通字段写入,若需保留文档 ID,应在source中读取_id并在 Sink 侧为其配置别名。
5.3 HTTPS/TLS 安全连接配置
场景一:关闭证书校验(适用于自签名证书的测试环境)
source { Easysearch { hosts = ["https://localhost:9200"] username = "admin" password = "admin" tls_verify_certificate = false } }场景二:关闭主机名校验(适用于 IP 直连或证书主机名不匹配的环境)
source { Easysearch { hosts = ["https://localhost:9200"] username = "admin" password = "admin" tls_verify_hostname = false } }场景三:启用证书校验并加载密钥库(生产环境推荐,使用 Easysearch 自带的http.p12证书)
source { Easysearch { hosts = ["https://localhost:9200"] username = "admin" password = "admin" tls_keystore_path = "${your Easysearch home}/config/certs/http.p12" tls_keystore_password = "${your password}" } }生产环境建议保持tls_verify_certificate与tls_verify_hostname均为true(默认值),并通过tls_keystore_path/tls_truststore_path配置可信证书链,避免关闭校验带来的中间人攻击风险。
六、运行与验证
6.1 运行前提
- 已安装 SeaTunnel 发行版(当前仓库为源码工程,正式使用请使用发布版安装包);
- 已安装或自行构建
connector-easysearch插件,并在 plugin_config 中启用该连接器; - Easysearch 集群可访问,且(如开启安全认证)具备读取目标索引的权限。
6.2 快速验证
将上述任一配置保存为easysearch_source.conf,使用 SeaTunnel 命令行提交(以 Zeta 引擎为例):
bin/seatunnel.sh --config easysearch_source.conf -e local6.3 单元测试验证
仓库内置了连接器的单元测试可参考:
- EasysearchSourceTest.java:验证「未配置 source 时依据 Easysearch 字段类型自动推导 SeaTunnel RowType」的逻辑;
- EasysearchSourceSplitEnumeratorTest.java:验证分片枚举器的分片分配行为;
- EasysearchFactoryTest.java:验证插件工厂注册与选项定义。
七、常见问题与使用建议
source和schema为什么互斥?二者都用于定义「读哪些字段、如何转换」:source只声明字段名(类型由索引 Mapping 推导或回退为 text),schema显式声明字段名与类型。同时配置会产生歧义,连接器会拒绝;都省略则自动采用索引的全部映射字段。通配索引的读取顺序:按索引文档数升序读取,便于先快速产出小索引的数据;若需严格控制顺序,建议在
index中精确指定单个索引。_id如何携带?在source中显式加入"_id"字段即可随行数据输出;但 Easysearch 不允许_id作为普通字段写入,写入目标索引时必须为其指定别名。大批量读取的调优方向:适当调大
scroll_size(如 500~1000)减少请求往返;确保scroll_time大于单批数据全量拉取的耗时;如数据源字段较多,优先用source裁剪只读所需列,减少网络与反序列化开销。HTTPS 报错排查:证书校验失败时优先检查
tls_keystore_path/tls_truststore_path的文件路径是否对 SeaTunnel 运行用户可读、密码是否正确;测试环境可临时关闭tls_verify_certificate/tls_verify_hostname定位问题,生产环境务必恢复校验。Easysearch 与 Elasticsearch 的区别:Easysearch 是 INFINI 发布的、与 Elasticsearch 协议兼容的搜索产品,连接器通过其 HTTP 协议与 Scroll API 通信,因此
hosts、index、query等配置在形态上与 Elasticsearch 类连接器相近,但本连接器仅面向 Easysearch 官方版本验证。
八、版本演进
连接器持续演进,主要变更记录(详见 connector-easysearch changelog):
- 2.3.5:新增对 INFINI Easysearch 的支持;
- 2.3.10:重构连接器通用选项、修复 SourceSplitEnumerator 日志名称错误;
- 2.3.11:支持 schema_save_mode / data_save_mode,补充 Source/Sink 状态类 serialVersionUID。
以上版本信息与当前仓库源码对应,具体行为以你实际使用的 SeaTunnel 版本为准。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考