
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、CatalogEasysearchCatalog三大部分本文聚焦其中的 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}?hindex,docsCountformatjson接口见 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_timestring否1mEasysearch 保持 Scroll 上下文存活的时间scroll_sizeint否100每次 Scroll 请求返回的最大记录数tls_verify_certificateboolean否true是否校验 HTTPS 证书tls_verify_hostnameboolean否true是否校验 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 端。当source与schema均未配置时连接器会调用/{index}/_mappings接口见 EasysearchClient.java获取索引 Mapping自动使用索引中的全部映射字段。3.5 schema [config]可选与 source 二选一数据的结构定义包含字段名与字段类型用于让 SeaTunnel 按显式类型定义转换所选字段。更多细节参见 Schema Feature。schema与source互斥。当两者都省略时连接器自动读取 Easysearch 字段映射并使用索引中所有映射字段。从DefaultSeaTunnelRowDeserializer的源码DefaultSeaTunnelRowDeserializer.java看schema支持的能力远不止标量类型Map 类型如mapstring, tinyintArray 类型如arraytinyint元素会按元素类型递归转换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]可选默认 1mEasysearch 为 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 optionsSource 插件通用参数详见 Source Common Options。四、数据类型映射Easysearch Source 内置了 Easysearch 字段类型到 SeaTunnel 类型的映射表源码实现于 EzsTypeMappingSeaTunnelType.javaEasysearch 数据类型SeaTunnel 数据类型STRING / KEYWORD / TEXTSTRINGBOOLEANBOOLEANBYTEBYTESHORTSHORTINTEGERINTLONGLONGFLOAT / HALF_FLOATFLOATDOUBLEDOUBLEDATELOCAL_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 mapstring, tinyint c_array arraytinyint 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 推导或回退为 textschema显式声明字段名与类型。同时配置会产生歧义连接器会拒绝都省略则自动采用索引的全部映射字段。通配索引的读取顺序按索引文档数升序读取便于先快速产出小索引的数据若需严格控制顺序建议在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 changelog2.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),仅供参考