ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SeaTunnel Elasticsearch Source 连接器实战指南:多索引同步、认证、SCROLL/PIT 分页与运行时字段

SeaTunnel Elasticsearch Source 连接器实战指南:多索引同步、认证、SCROLL/PIT 分页与运行时字段 SeaTunnel Elasticsearch Source 连接器实战指南多索引同步、认证、SCROLL/PIT 分页与运行时字段【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 Elasticsearch Source 连接器插件名Elasticsearch用于从 Elasticsearch 集群批量读取数据支持 ES 2.x 到 8.x 之间各主要版本可作为批处理管道的入口将索引数据同步到任意目标端。读完本文你将掌握该连接器的全部配置参数、三种查询方式DSL / SQL、两种分页 APISCROLL / PIT的选型原则以及多索引并行同步、切片加速、TLS 认证与运行时字段等生产级用法。简介与能力边界Elasticsearch Source 连接器在 connector-elasticsearch 模块 中实现是一个有界BOUNDED批处理数据源。从源码看ElasticsearchSource.java 实现了SeaTunnelSource、SupportParallelism、SupportColumnProjection三个核心接口这与下表所列特性一一对应。特性支持情况批处理✔ 支持流处理✘ 不支持精准一次✘ 不支持列投影字段裁剪✔ 支持并行度✔ 支持用户自定义分片✘ 不支持但支持基于slice_max的服务端切片并行版本兼容说明连接器支持读取 Elasticsearch 2.x 至 8.x 之间各版本的数据。个别高级能力对版本有硬性要求例如 SCROLL 切片sliced scroll要求 ES ≥ 5.0PIT 要求 ES ≥ 7.10runtime_fields要求 ES ≥ 7.11详见下文对应小节。读取流程从 Split 枚举到数据行输出理解底层执行模型有助于正确配置并行度和切片参数。该连接器严格遵循 SeaTunnel 的 Source/Split 模型核心链路如下ElasticsearchSource 解析配置读取index或index_list生成ElasticsearchConfig列表如果配置了source字段会调用EsRestClient.getFieldTypeMapping拉取索引映射mapping推导 SeaTunnel 字段类型若为 SQL 模式则调用getSqlMapping从 SQL 结果推导类型。ElasticsearchSourceSplitEnumerator 枚举 Split通过getIndexDocsCount获取每个索引的文档数过滤掉空索引后为每个(索引, 切片)组合生成一个 Splitslice_max 1时 splitId 形如index#sliceId再按assignCount % readerCount轮询round-robin分配给各并行子任务。ElasticsearchSourceReader 消费数据按search_type分支执行——SQL 走 X-Pack SQL 游标分页DSL 模式再按search_api_type决定走 Scroll API 还是 PIT API逐批拉取并用DefaultSeaTunnelRowDeserializer反序列化为SeaTunnelRow交给下游。由于getBoundedness()返回BOUNDED读取完毕会发送signalNoMoreElement结束任务因此该连接器天然适合一次性全量/增量窗口同步不适合持续监听写入。配置参数总览所有参数均在 ElasticsearchSourceOptions.java 与 ElasticsearchBaseOptions.java 中声明含默认值汇总如下参数名称类型是否必须默认值或说明hostsarray是Elasticsearch 集群 HTTP 地址格式host:port可配多个auth_typestring否basicusernamestring否x-pack 用户名passwordstring否x-pack 密码auth.api_key_idstring否API Key 的 IDauth.api_keystring否API Key 的密钥auth.api_key_encodedstring否Base64 编码的 API Keybase64(id:key)indexstring否索引名支持*通配未配置index_list时必须配置index_listarray否多索引同步任务定义sourcearray否要读取的字段列表不配置则从索引映射自动获取queryjson否{match_all: {}}search_typeenum否查询类型DSL或SQL默认DSLsearch_api_typeenum否分页 API 类型SCROLL或PIT默认SCROLLsql_querystring否SQL 查询语句search_type SQL时必填scroll_timestring否1m搜索上下文存活时长scroll_sizeint否100每次滚动返回的最大文档数tls_verify_certificateboolean否truetls_verify_hostnameboolean否truearray_columnmap否声明数组字段类型ES 本身无数组类型tls_keystore_pathstring否PEM 或 JKS 密钥库路径tls_keystore_passwordstring否密钥库密码tls_truststore_pathstring否PEM 或 JKS 信任库路径tls_truststore_passwordstring否信任库密码pit_keep_alivelong否60000毫秒即 1 分钟pit_batch_sizeint否100slice_maxint否11 时启用切片并行读取SCROLL 需 ES≥5.0PIT 需 ES≥7.10runtime_fieldsarray否查询时动态计算的字段ES 7.11common-options-否Source 插件通用参数基础连接与认证hosts [array]必填Elasticsearch 集群的 HTTP 地址数组格式为host:port支持一次配置多个节点实现负载均衡与故障转移。例如[host1:9200, host2:9200]。从 EsRestClient.java 的实现看hosts 会被逐个解析为HttpHost构建底层RestClient并为每个请求设置 10 秒连接请求超时与 5 分钟 Socket 超时。auth_type [enum]指定认证方式由 AuthenticationProviderFactory 根据配置分发到对应 Provider。支持basic默认用户名 密码的 HTTP 基本认证对应 x-pack 安全api_keyAPI Key 的 ID 密钥认证api_key_encodedBase64 编码后的 API Key 认证。未指定时默认使用basic以兼容旧版本配置。基本认证basic参数类型说明usernamestring基本认证用户名x-pack 用户名passwordstring基本认证密码x-pack 密码source { Elasticsearch { hosts [https://localhost:9200] auth_type basic username elastic password your_password index my_index } }API Key 认证参数类型说明auth.api_key_idstringElasticsearch 生成的 API Key IDauth.api_keystringElasticsearch 生成的 API Key 密钥auth.api_key_encodedstringbase64(id:api_key)形式的编码 Key可替代单独提供 ID 与 key注意auth.api_key_idauth.api_key与auth.api_key_encoded只能二选一同时配置时工厂校验会将其视为非法组合。示例分开配置 ID 和 keysource { Elasticsearch { hosts [https://localhost:9200] auth_type api_key auth.api_key_id your_api_key_id auth.api_key your_api_key_secret index my_index } }示例使用编码 keysource { Elasticsearch { hosts [https://localhost:9200] auth_type api_key_encoded auth.api_key_encoded eW91cl9hcGlfa2V5X2lkOnlvdXJfYXBpX2tleV9zZWNyZXQ index my_index } }索引与字段投影配置index [string]Elasticsearch 索引名称支持*通配符。例如存在索引index1、index2可指定index*同时读取两个索引的数据。index与index_list至少配置一个单索引或索引通配场景使用index不同索引需要单独配置query、source、schema或分页参数时使用index_list。行为提示从 ElasticsearchSource.java 的构造逻辑看若index与index_list同时出现会打印告警并只让index_list生效。source [array]要读取的索引字段列表即列投影。你可以通过指定字段_id来获取文档 ID如果要将_id写入其他索引由于 Elasticsearch 的限制需要为_id指定一个别名。如果未配置source连接器会通过EsRestClient.getFieldTypeMapping自动从索引映射中获取全部字段及其类型。字段顺序即输出 SeaTunnel Row 的字段顺序。array_column [map]由于 Elasticsearch 中没有数组类型映射推导无法识别数组字段因此需要通过该参数显式声明。假设tags和phones是数组字段array_column {tags arraystring, phones arraystring}在 ElasticsearchSource.java 中array_column中命中的字段会直接按其声明的 SeaTunnel 类型如arraystring、arraytinyint构建物理列其余字段才走 ES 类型到 SeaTunnel 类型的自动转换ElasticSearchTypeConverter。schema已废弃早期版本使用schema显式声明字段类型但源码中已明确打印警告The schema config in ElasticSearch source/sink is deprecated, please use source config instead!。当前推荐做法是只配置source字段列表配合可选的array_column让连接器根据索引映射自动推导类型。查询与分页query [json]Elasticsearch 原生查询语句DSL用于控制读取哪些文档。不配置时默认值为{match_all: {}}即全量读取。例如{range:{firstPacket:{gte:1669225429990,lte:1669225429990}}}search_type [enum] 与 sql_query [string]DSL默认使用 Elasticsearch Query DSL 查询SQL使用 X-Pack SQL 查询此时必须配置sql_query。SQL 模式通过EsRestClient.searchBySql发起/_sql请求并基于服务端返回的 cursor 游标持续翻页结束后调用closeSqlCursor释放资源。source { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index st_index_sql sql_query select * from st_index_sql where c_int10 and c_int20 search_type sql } }注意SQL 查询不支持 map 和 array 类型字段同时 SQL 查询不支持切片slice_max配置会被忽略源码中会打印告警并强制置回 1。scroll_time [string] 与 scroll_size [int]底层使用滚动查询Scroll API拉取数据因此需要scroll_time控制搜索上下文search context在 Elasticsearch 侧存活的时长默认1m1 分钟。若任务处理较慢导致上下文过期会触发SearchContextMissingException此时应适当调大。scroll_size每次滚动请求返回的最大文档数默认100。读取完成后连接器会调用clearScroll主动清理 scroll ID。search_api_type [enum]SCROLL 与 PIT 对比参数值分页机制版本要求特点SCROLL默认Scroll APIES 2.x 起即可用实现简单、兼容性最广但滚动过程中数据变更可能影响结果一致性PITPoint in Time API search_afterES ≥ 7.10基于共享快照跨切片/跨页一致性更好需要维护 PIT IDpit_keep_alive[long]PIT 保持活动的时间量单位毫秒默认600001 分钟。pit_batch_size[int]每次 PIT 搜索请求返回的最大文档数默认100语义同scroll_size。PIT 模式下 ElasticsearchSourceReader 会先创建 PIT随后用search_after游标持续翻页直至hasMore为 false最后在finally中删除 PIT 释放资源。slice_max [int]切片并行读取将单个索引拆分为多个切片slice并行读取仅对 SCROLL/PIT 生效配置 1 时启用。版本要求SCROLL 切片sliced scroll需要 Elasticsearch 5.0 及以上版本PIT 切片需要 Elasticsearch 7.10 及以上版本PIT 于 7.10.0 引入。取舍说明切片能显著提升吞吐但可能降低跨切片的数据一致性。对一致性要求高的场景建议使用 PIT共享快照或将slice_max 1对追加写或写入较少的场景开启切片通常可以接受。当search_type SQL时Elasticsearch SQL 查询不支持切片slice_max会被忽略。与并行度的配合从 ElasticsearchSourceSplitEnumerator.java 的枚举逻辑可见slice_max决定每个索引生成的 Split 数量sliceId从 0 到sliceMax-1而parallelism通用选项决定参与消费的 Reader 子任务数所有 Split 会按轮询方式均匀分配到各 Reader。因此提升吞吐通常需要同时调大这两者。runtime_fields查询期动态计算字段runtime_fields允许在查询时动态计算字段值而无需重建索引适合临时分析、字段试验与低频查询Elasticsearch 7.11。每个 runtime field 需要包含name字段名type数据类型支持boolean、date、double、geo_point、ip、keyword、longscriptPainless 脚本用于计算字段值script_lang可选脚本语言默认painlessscript_params可选脚本参数在源码 ElasticsearchSource.java 的parseRuntimeFields方法中该配置会被转换为 Elasticsearch 的runtime_mappings结构{字段名: {type, script: {source, lang?, params?}}}随查询请求一并下发。runtime_fields [ { name day_of_week type keyword script emit(doc[timestamp].value.dayOfWeekEnum.toString()) }, { name total_price type double script emit(doc[quantity].value * doc[price].value) } ]性能与限制运行时字段在查询阶段计算数据量大时会影响查询性能适合临时分析、字段试验与低频查询场景需要 Elasticsearch 7.11 及以上版本计算出的字段需加入source列表并建议在schema中声明类型后才可被下游使用。TLS / SSL 加密连接配置连接器基于 ElasticsearchBaseOptions.java 提供完整的 HTTPS 支持参数类型默认值说明tls_verify_certificatebooleantrue是否启用 HTTPS 端点的证书验证tls_verify_hostnamebooleantrue是否启用 HTTPS 端点的主机名验证tls_keystore_pathstring-PEM 或 JKS 密钥库路径文件必须对运行 SeaTunnel 的操作系统用户可读tls_keystore_passwordstring-密钥库的密钥密码tls_truststore_pathstring-PEM 或 JKS 信任库路径文件必须对运行 SeaTunnel 的操作系统用户可读tls_truststore_passwordstring-信任库的密钥密码典型使用方式是 hosts 使用https://前缀并配合相应证书策略完整示例见下文「使用案例」的案例三至案例五。多索引同步index_listindex_list用于定义多索引同步任务。它是一个数组每个元素包含单表同步所需的完整参数如index、query、source/schema、scroll_size和scroll_time等。建议不要将index_list和query配置在同一层级即query应放在index_list的每个条目内。注意源码注释披露的限制index_list内部每个条目的配置不会经过工厂的OptionRule校验该校验只作用于顶层配置。如果某个条目缺少index或设置为search_type SQL却缺少sql_query会在运行时解析阶段才报错。因此多索引配置时务必逐个核对条目完整性。common options通用选项Source 插件常用参数完整说明见 Source 常用选项名称类型必填说明plugin_outputString否将本插件数据注册为可被其他插件直接访问的数据集/临时表旧名result_table_name已过时parallelismInt否覆盖环境中的并行度设置未指定时使用环境默认值metadata_datasource_idString否从元数据中心获取连接配置的数据源 ID重要提示作业中使用plugin_output时下游插件必须设置plugin_input才能消费该数据集。使用案例案例一通配索引 字段投影 数组字段 范围查询从满足seatunnel-*匹配的索引中按 query 读取数据查询只返回文档的id、name、age、tags、phones字段其中tags、phones通过array_column声明为数组类型_id用于取回文档 ID。Elasticsearch { hosts [localhost:9200] index seatunnel-* array_column {tags arraystring,phones arraystring} source [_id,name,age,tags,phones] query {range:{firstPacket:{gte:1669225429990,lte:1669225429990}}} }案例二多索引同步index_list演示从read_index1和read_index2读取不同的数据。read_index1使用source指定字段、array_column声明数组字段并带 range 条件read_index2使用match_all全量读取另一组字段。两个来源经同一管道由 Elasticsearch Sink 统一写出示例 sink 目标索引为multi_source_write_test_indexindex_type st采用CREATE_SCHEMA_WHEN_NOT_EXISTAPPEND_DATA的保存模式。source { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index_list [ { index read_index1 query {range: {c_int: {gte: 10, lte: 20}}} source [ c_map, c_array, c_string, c_boolean, c_tinyint, c_smallint, c_bigint, c_float, c_double, c_decimal, c_bytes, c_int, c_date, c_timestamp ] array_column { c_array arraytinyint } } { index read_index2 query {match_all: {}} source [ c_int2, c_date2, c_null ] } ] } } transform { } sink { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index multi_source_write_test_index index_type st schema_save_modeCREATE_SCHEMA_WHEN_NOT_EXIST data_save_modeAPPEND_DATA } }案例三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 } }案例五SSL启用证书验证通过密钥库此处以 ES 安装目录下的http.p12为例建立双向信任生产环境推荐方式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} } }案例六SQL 方式查询使用 X-Pack SQL 语法查询。注意SQL 查询不支持 map 和数组类型字段。source { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index st_index_sql sql_query select * from st_index_sql where c_int10 and c_int20 search_type sql } }案例七PIT 方式滚动查询使用 DSL 查询 Point in Time 分页pit_keep_alive为 60000 毫秒1 分钟每次拉取 100 条source { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index st_index query {range: {c_int: {gte: 10, lte: 20}}} # 使用 DSL 查询和 PIT API search_type DSL search_api_type PIT pit_keep_alive 60000 # 1 minute in milliseconds pit_batch_size 100 } }案例八Runtime Fields 查询期计算字段在查询时计算字段值而无需重建索引。定义 4 个运行时字段含条件分支与脚本参数并声明输出字段与 schema 后写入 Consolesource { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index sales_data # 定义运行时字段 runtime_fields [ { name total_amount type double script emit(doc[quantity].value * doc[price].value) }, { name day_of_week type keyword script emit(doc[order_date].value.dayOfWeekEnum.getDisplayName(TextStyle.FULL, Locale.ROOT)) }, { name order_category type keyword script double amount doc[quantity].value * doc[price].value; if (amount 1000) { emit(high_value); } else if (amount 100) { emit(medium_value); } else { emit(low_value); } }, { name price_with_tax type double script emit(doc[price].value * (1 params.tax_rate)) script_params { tax_rate 0.13 } } ] source [ product_id, quantity, price, order_date, total_amount, day_of_week, order_category, price_with_tax ] schema { fields { product_id string quantity int price double order_date timestamp total_amount double day_of_week string order_category string price_with_tax double } } } } sink { Console { } }说明虽然schema已标记为废弃但示例中保留它是为了给运行时字段声明明确的输出类型纯source模式无法推导运行时字段类型因此在包含runtime_fields的场景中仍需以schema补齐类型信息。案例九PIT slicing 并行读取在 PIT 分页基础上开启slice_max 2将st_index拆分为 2 个切片并行拉取提升吞吐source { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index st_index query {range: {c_int: {gte: 10, lte: 20}}} search_type DSL search_api_type PIT pit_keep_alive 60000 pit_batch_size 100 # 开启切片并行读取 slice_max 2 } }版本演进与变更记录该连接器自 SeaTunnel 2.2.0-beta 引入 ES Sink、2.3.0 正式支持 Source 以来持续演进完整的逐版本变更见 connector-elasticsearch 变更日志与本 Source 相关的重要里程碑包括2.3.8支持多表源multi-table source特性即index_list2.3.10支持 Elasticsearch SQL Sourcesearch_type SQL2.3.11支持 PIT 分页search_api_type PIT2.3.12新增 API Key 认证支持auth_type api_key/api_key_encoded。配置建议小结连接生产环境优先 HTTPS 证书校验案例五测试环境可临时关闭校验案例三/四启用安全认证的集群按需选择 basic 或 API Key。查询默认 DSL query即可满足绝大多数场景需要复杂聚合/投影计算时改用 SQL但要避开 map/array 字段并接受无法切片的限制。分页ES ≥ 7.10 且对一致性有要求时优先 PIT共享快照老版本集群用 SCROLL 并将scroll_time调大到足以覆盖最慢批次的处理耗时。吞吐大索引开启slice_max 1配合调大parallelism数据持续追加、可容忍近似一致性的场景收益最大。多索引结构不同、条件不同的索引使用index_list逐条配置避免把query放在顶层与index_list混用。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表