数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本篇技术指南围绕 Apache SeaTunnel 的Elasticsearch Source 连接器展开介绍如何在 SeaTunnel 批处理作业中从 Elasticsearch含 OpenSearch 兼容场景2.x 至 8.x 版本读取数据完整覆盖hosts、index、query、scroll_time、scroll_size、source、array_column及全部 TLS 相关参数的配置方法并结合本仓库的源码实现讲解底层 Scroll 分页、Split 分配与类型映射机制。读完本文你将能够独立编写一个可运行的 Elasticsearch → SeaTunnel 数据读取作业并理解其内部执行原理便于排查问题与二次开发。连接器概述Elasticsearch Source 连接器用于从 Elasticsearch 集群读取数据官方文档位于 docs/en/connector-v2/source/Elasticsearch.md其插件实现位于 connector-elasticsearch 模块。连接器支持 Elasticsearch2.x 至 8.x版本并兼容基于 HTTPS 协议的 OpenSearch 场景见文末 Changelog 中 [PR 3997] 的相关描述。支持的功能特性根据文档中的特性矩阵特性定义详见 Connector V2 Features 说明特性支持情况说明batch批处理✅数据有界读取完所有数据后作业自动结束stream流处理❌不支持持续读取新增数据exactly-once精确一次❌未提供基于 offset 的断点续读保证column projection列投影✅通过source参数只读取指定字段parallelism并行度⚠️文档特性矩阵未勾选但源码实现了SupportParallelism接口详见下文并行读取机制support user-defined split自定义分片❌分片规则由连接器内部决定配置选项全解析以下是文档给出的完整选项表名称类型是否必填默认值hostsarray是-usernamestring否-passwordstring否-indexstring是-sourcearray否-queryjson否{match_all: {}}scroll_timestring否1mscroll_sizeint否100tls_verify_certificateboolean否truetls_verify_hostnamesboolean否truearray_columnmap否-tls_keystore_pathstring否-tls_keystore_passwordstring否-tls_truststore_pathstring否-tls_truststore_passwordstring否-common-options-否-注意表格中写作tls_verify_hostnames但源码 EsClusterConnectionConfig.java 中实际定义的配置键为tls_verify_hostname官方示例也使用tls_verify_hostname配置时应以此为准。各选项定义与默认值均可在 SourceConfig.java 与 EsClusterConnectionConfig.java 中找到对应实现。hosts必填arrayElasticsearch 集群的 HTTP 地址列表格式为host:port可配置多个节点例如[host1:9200, host2:9200]。底层通过 Elasticsearch Low-Level REST Client 构建HttpHost数组建立连接详见 EsRestClient.java 中getRestClientBuilder方法。当启用 HTTPS 时地址应写为https://host:port如https://localhost:9200。username / password可选stringElasticsearch x-pack 安全认证的用户名与密码。源码中通过BasicCredentialsProvider设置UsernamePasswordCredentials注入 HttpClient用于访问受安全保护的集群。index必填string要读取的 Elasticsearch 索引名称支持*通配符模糊匹配例如seatunnel-*可一次读取多个索引。读取前连接器会调用/_cat/indices/{index}?hindex,docsCountformatjson接口获取每个匹配索引的文档数见 EsRestClient.java 的getIndexDocsCount方法并且只读取文档数大于 0 的索引。source可选array要读取的索引字段列表即列投影。例如source [_id, name, age]。通过指定_id可以获取文档 ID如果后续要将_id写入其他索引由于 Elasticsearch 的限制需要为_id指定别名alias。如果不配置source连接器会自动从索引的 mapping 中获取全部字段见下方字段类型自动推断小节。array_column可选mapES 中没有数组类型索引因此需要手动为数组字段指定类型格式形如{c_array arraytinyint}。文档表格将其标注为 map 类型小节标题写作[array]属于笔误源码 SourceConfig.java 中定义类型为MapString, String默认值为空 Map。query可选jsonElasticsearch 查询 DSL用于控制读取的数据范围。默认值为{match_all: {}}即读取全部文档。源码中该配置直接作为请求体中的query字段传给 ES见searchByScroll方法。典型用法如{range:{firstPacket:{gte:1669225429990,lte:1669225429990}}}按时间范围过滤。注意此处无法使用聚合类 DSL滚动读取的是hits.hits中的原始文档。scroll_time可选stringElasticsearch 为 Scroll 请求保持搜索上下文search context的存活时间默认1m。每次滚动请求都会以该时间刷新上下文有效期。需要根据数据处理速度合理设置过短可能导致大批量数据读取中途上下文过期报错过长会占用服务端资源。scroll_size可选int每次 Elasticsearch Scroll 请求返回的最大命中数默认100。该值影响单批拉取的数据量与网络往返次数可结合实际吞吐需求调整。TLS / SSL 相关参数tls_verify_certificateboolean默认true是否校验 HTTPS 端点的证书。设为false时使用 Trust-All 策略跳过证书校验源码中构造TrustAllStrategy的 SSLContext。tls_verify_hostnameboolean默认true是否校验 HTTPS 端点的主机名。设为false时使用NoopHostnameVerifier跳过主机名校验。tls_keystore_pathstringPEM 或 JKS 格式密钥库key store路径运行 SeaTunnel 的操作系统用户必须可读该文件。tls_keystore_passwordstring上述密钥库的密钥密码。tls_truststore_pathstringPEM 或 JKS 格式信任库trust store路径同样要求运行用户可读。tls_truststore_passwordstring上述信任库的密钥密码。在 EsRestClient.java 的createInstance中可以看到仅当tls_verify_certificate为true时才会读取 keystore/truststore 相关配置并通过SSLUtils.buildSSLContext构建 SSLContext证书校验与主机名校验是两个独立的开关可分别关闭。common-options可选Source 插件的公共参数包括result_table_name注册结果表供下游插件通过source_table_name引用与parallelism覆盖 env 中的并行度等详见 Source Common Options。数据读取核心原理Scroll 分页 Split 机制Elasticsearch Source 是批处理BOUNDED源读取过程由ElasticsearchSourceSplitEnumerator分片枚举器与ElasticsearchSourceReader读取器协同完成属于典型的 SeaTunnel Source API 拆分-读取模型。1. Split 枚举阶段ElasticsearchSourceSplitEnumerator.java 中的getElasticsearchSplit方法完成以下动作通过/_cat/indices/{index}?hindex,docsCountformatjson获取所有匹配索引及其文档数过滤掉文档数为 0 的索引并按文档数升序排序为每个索引创建一个ElasticsearchSourceSplitsplitId 为索引名 hashCodesplit 中携带索引名、字段列表source、query、scroll_time、scroll_size 等信息依据currentParallelism()当前并行度通过(splitId.hashCode() Integer.MAX_VALUE) % numReaders将 split 哈希分配到各个 reader 任务。并行读取机制虽然特性矩阵未勾选 parallelism但从源码结构看ElasticsearchSource实现了SupportParallelism接口见 ElasticsearchSource.java且枚举器会按并行度将不同索引的 split 分配给不同 reader。因此当index匹配多个索引时可以通过 common-options 的parallelism参数让多个任务并行读取不同索引若只有一个索引则退化为单任务顺序读取。2. Reader 滚动读取阶段ElasticsearchSourceReader.java 的pollNext方法对每个 split 执行 Scroll 流程首次请求调用EsRestClient.searchByScroll向/{index}/_search?scroll{scrollTime}发起 POST请求体包含query、_source字段投影、sort: [_doc]按内部文档顺序排序避免排序开销以及size: scroll_size循环续滚只要上一批docs非空就携带scroll_id调用searchWithScrollId接口/_search/scroll继续拉取下一批直到返回空批次为止每批结果通过DefaultSeaTunnelRowDeserializer反序列化为SeaTunnelRow后交给Collector下发。ES 服务端返回的每个文档会被重组为一个 Map固定放入_index、_id字段_source中的各字段平铺进 Map字符串类型保留为文本其余以 JsonNode 形式保留见getDocsFromScrollResponse方法。3. 失败恢复Checkpoint 支持Reader 与 Enumerator 均实现了snapshotStateReader 快照未消费的 split 队列Enumerator 快照待分配 split 与shouldEnumerate标记。任务重启后可基于快照恢复读取进度这与exactly-once语义不同——Scroll 上下文本身是易失的快照仅保证 split 级别的断点续跑。字段类型推断与类型转换自动推断字段类型当未配置source时连接器调用EsRestClient.getFieldTypeMapping请求/{index}/_mappings接口解析索引 mapping 的properties将每个字段的 ES 类型映射为 SeaTunnel 类型通过 ElasticSearchTypeConverter 完成并自动以全部字段作为source列表构建CatalogTable。从 getFieldTypeMappingFromProperties 的实现可以看到对特殊类型的处理aggregate_metric_double携带metrics子选项alias解析path指向的真实字段类型dense_vector记录element_type默认floatdate/date_nanos记录format默认strict_date_optional_time_nanos||epoch_millis含properties的字段按object类型递归解析若在 mapping 中找不到某字段则回退为text类型并打印告警日志。值反序列化DefaultSeaTunnelRowDeserializer.java 负责将文档值按 SeaTunnel 类型转换字符串、布尔、byte/short/int/long/float/double直接解析date/time/timestamp先尝试按 epoch 毫秒时间戳解析否则支持yyyy-MM-dd HH、yyyy-MM-dd HH:mm:ss、yyyy-MM-dd HH:mm:ss.SSS、yyyyMMdd HH:mm:ss等十余种格式按字符串长度匹配格式化器并兼容T/Z前缀替换decimalnew BigDecimal(value)array将 JSON 数组字符串解析为 List 后逐个转换元素构造对应类型的数组map解析为MapString, String后按 key/value 类型递归转换嵌套 row递归解析子文档字段bytes按 Base64 解码值为字符串null或 JSONNullNode时置为 null。该反序列化逻辑与array_column配置配合使用对标记为数组的字段即使 ES 中没有数组类型也能正确还原为 SeaTunnel 的数组列。仓库测试 ElasticsearchSourceTest.java 验证了空source场景下字段类型推断的可用性。配置示例以下示例均来自官方文档可直接复制到 SeaTunnel 作业配置HOCON 格式中使用。完整作业还需env与sink部分例如配合 Console sink 将结果打印到控制台。示例一简单读取Elasticsearch { hosts [localhost:9200] index seatunnel-* source [_id,name,age] query {range:{firstPacket:{gte:1669225429990,lte:1669225429990}}} }该示例读取seatunnel-前缀的所有索引仅取_id、name、age三个字段并按firstPacket字段的时间范围毫秒时间戳过滤数据。示例二复杂类型显式 SchemaElasticsearch { hosts [elasticsearch:9200] index st_index 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 } } query {range:{firstPacket:{gte:1669225429990,lte:1669225429990}}} }该示例通过schema.fields显式声明各字段的 SeaTunnel 类型覆盖 map、array、decimal、bytes、date、timestamp 等复杂类型。需要注意源码 ElasticsearchSource.java 在检测到schema配置时会打印The schema config in ElasticSearch sink is deprecated, please use source config instead!的告警提示该用法已过时——新版本推荐使用source字段列表 array_column的方式声明字段官方文档仍保留此示例供兼容参考。示例三SSL关闭证书校验source { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_verify_certificate false } }适用于自签名证书场景跳过证书链校验但仍保留主机名校验。示例四SSL关闭主机名校验source { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_verify_hostname false } }适用于证书合法但主机名不匹配如访问 IP 而非域名的场景。示例五SSL启用证书校验使用密钥库source { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_keystore_path ${your elasticsearch home}/config/certs/http.p12 tls_keystore_password ${your password} } }这是生产环境推荐的安全配置指定 PEM/JKS 密钥库路径与密码完成双向 TLS 认证注意文件需对运行 SeaTunnel 的操作系统用户可读。注意事项与使用建议版本范围连接器支持 Elasticsearch 2.x 至 8.x。由于底层使用 Low-Level REST Client 直接调用_search、_scroll、_mappings、_cat/indices等 HTTP 接口不依赖 ES 高版本 SDK 的强类型绑定因此跨版本兼容性较好但在实际部署前仍建议以目标集群版本做一次连通性验证。查询字段query中无法使用聚合、脚本字段script_fields等会改变返回结构的功能读取范围控制请使用 Query DSL 的过滤能力。Scroll 上下文开销scroll_time期间 ES 会为每个 Scroll 保留搜索上下文并占用堆内存请勿设置过大的scroll_size与过长的scroll_time读取结束后连接器不会显式调用_clear_scroll依赖上下文自动过期。数组字段必须显式声明ES 索引没有数组类型读取数组字段时必须通过array_column指定否则会被按默认类型解析。特性边界该连接器不支持流模式与 exactly-once若需增量同步可考虑结合 CDC 类连接器参见 connector-cdc 模块如需将数据写回 Elasticsearch参见 Elasticsearch Sink 文档。变更记录根据文档 Changelog该连接器随 SeaTunnel 演进新增了以下能力next version 条目新增 Elasticsearch Source Connector支持 HTTPS 协议并兼容 OpenSearch支持 Query DSL。对应的功能均可在当前仓库源码中找到实现HTTPS/TLS 支持位于 EsRestClient.java 的连接构建逻辑Query DSL 支持位于 SourceConfig.java 的QUERY选项定义及searchByScroll的请求拼装逻辑中。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Easysearch Source 连接器详解从 Scroll 读取原理到 TLS 安全配置SeaTunnel Easysearch Source 连接器详解从 Scroll 读取原理到 TLS 安全配置 本篇技术指南围绕 SeaTunnelApa数据工程大数据批处理流处理SeaTunnel ClickHouse Source Connector 全解析从参数配置到并行读取原理SeaTunnel ClickHouse Source Connector 全解析从参数配置到并行读取原理 ClickHouse 以其极致的列式存储与向量化查数据集成ETL大数据批处理流处理变更数据捕获使用 SeaTunnel Hbase Source Connector 读取 Apache Hbase 数据完整配置指南与原理剖析使用 SeaTunnel Hbase Source Connector 读取 Apache Hbase 数据完整配置指南与原理剖析 本文是 SeaTunnel数据工程大数据批处理流处理上一篇React Native Background Geolocation调试与排错10个常见问题终极解决方案下一篇3步解锁yuzu模拟器高帧率从技术瓶颈到丝滑体验的完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考