SeaTunnel Easysearch Source Connector 完全指南:从 INFINI Easysearch 批量读取数据
2026/9/17 14:05:19 网站建设 项目流程

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,即这是一个有界(批式)数据源;同时该类实现了SupportParallelismSupportColumnProjection接口——需要注意的是,虽然连接器声明支持并行接口,但每个并行子任务都会对同一批索引发起独立的 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:

名称类型是否必填默认值说明
hostsarray-Easysearch HTTP 地址列表
indexstring-Easysearch 索引名,支持seatunnel-*通配匹配
usernamestring-安全认证用户名
passwordstring-安全认证密码
sourcearray-需要读取的字段列表;与schema二选一
schemaconfig-用于读取与转换字段的 SeaTunnel Schema;与source二选一
queryjson{"match_all":{}}Easysearch DSL 查询,用于过滤记录
scroll_timestring1mEasysearch 保持 Scroll 上下文存活的时间
scroll_sizeint100每次 Scroll 请求返回的最大记录数
tls_verify_certificatebooleantrue是否校验 HTTPS 证书
tls_verify_hostnamebooleantrue是否校验 HTTPS 主机名
tls_keystore_pathstring-PEM 或 JKS 密钥库路径
tls_keystore_passwordstring-密钥库密码
tls_truststore_pathstring-PEM 或 JKS 信任库路径
tls_truststore_passwordstring-信任库密码
common-optionsconfig-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 端)。

sourceschema均未配置时,连接器会调用/{index}/_mappings接口(见 EasysearchClient.java)获取索引 Mapping,自动使用索引中的全部映射字段。

3.5 schema [config](可选,与 source 二选一)

数据的结构定义,包含字段名与字段类型,用于让 SeaTunnel 按显式类型定义转换所选字段。更多细节参见 Schema Feature。

schemasource互斥。当两者都省略时,连接器自动读取 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:ssyyyy-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}创建上下文,后续每次翻页都会携带该参数续期。该值需大于单次拉取全部数据的耗时,否则上下文过期会导致读取中断;数据量大时建议调大(如2m5m)。

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 / TEXTSTRING
BOOLEANBOOLEAN
BYTEBYTE
SHORTSHORT
INTEGERINT
LONGLONG
FLOAT / HALF_FLOATFLOAT
DOUBLEDOUBLE
DATELOCAL_DATE_TIME_TYPE

几点源码层面的补充:

  1. 映射表不区分大小写场景有限:映射键为小写(stringkeywordtextintegerhalf_float等),Mapping 中获取的类型会被直接用于查表;
  2. 源码中还额外支持了binary类型映射到 STRING(文档表中未单列);
  3. 当 Easysearch 类型不在映射表内时,会抛出EasysearchConnectorException(错误码EZS_FIELD_TYPE_NOT_SUPPORT),例如nestedobject等复合类型若直接参与自动 Mapping 转换可能报错,此时建议用schema显式声明字段类型;
  4. 上述映射主要用于「未配置 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-*索引,读取每个文档的_idnameage三个字段,且仅保留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_indexc_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_certificatetls_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 local

6.3 单元测试验证

仓库内置了连接器的单元测试可参考:

  • EasysearchSourceTest.java:验证「未配置 source 时依据 Easysearch 字段类型自动推导 SeaTunnel RowType」的逻辑;
  • EasysearchSourceSplitEnumeratorTest.java:验证分片枚举器的分片分配行为;
  • EasysearchFactoryTest.java:验证插件工厂注册与选项定义。

七、常见问题与使用建议

  1. sourceschema为什么互斥?二者都用于定义「读哪些字段、如何转换」:source只声明字段名(类型由索引 Mapping 推导或回退为 text),schema显式声明字段名与类型。同时配置会产生歧义,连接器会拒绝;都省略则自动采用索引的全部映射字段。

  2. 通配索引的读取顺序:按索引文档数升序读取,便于先快速产出小索引的数据;若需严格控制顺序,建议在index中精确指定单个索引。

  3. _id如何携带?source中显式加入"_id"字段即可随行数据输出;但 Easysearch 不允许_id作为普通字段写入,写入目标索引时必须为其指定别名。

  4. 大批量读取的调优方向:适当调大scroll_size(如 500~1000)减少请求往返;确保scroll_time大于单批数据全量拉取的耗时;如数据源字段较多,优先用source裁剪只读所需列,减少网络与反序列化开销。

  5. HTTPS 报错排查:证书校验失败时优先检查tls_keystore_path/tls_truststore_path的文件路径是否对 SeaTunnel 运行用户可读、密码是否正确;测试环境可临时关闭tls_verify_certificate/tls_verify_hostname定位问题,生产环境务必恢复校验。

  6. Easysearch 与 Elasticsearch 的区别:Easysearch 是 INFINI 发布的、与 Elasticsearch 协议兼容的搜索产品,连接器通过其 HTTP 协议与 Scroll API 通信,因此hostsindexquery等配置在形态上与 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询