拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

SeaTunnel Elasticsearch Sink 连接器完全指南:参数详解、CDC 写入与源码级实现剖析

SeaTunnel Elasticsearch Sink 连接器完全指南:参数详解、CDC 写入与源码级实现剖析 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载SeaTunnel 内置的ElasticsearchSink 插件用于将 SeaTunnel 管道中的行数据SeaTunnelRow批量写入 Elasticsearch 集群支持动态索引名、主键文档_id生成、CDCChange Data Capture事件INSERT/UPDATE/DELETE以及 HTTPS/TLS 安全连接兼容 Elasticsearch 2.x ~ 8.x。读完本文你将能够从零配置一个可运行的 Elasticsearch Sink 作业理解每个参数的底层影响并掌握其 Bulk 批量写入、重试、序列化与 SaveMode 的实现原理。概述与能力边界ElasticsearchSink 插件位于 connector-elasticsearch 模块通过 plugin-mapping.properties 中的seatunnel.sink.Elasticsearch connector-elasticsearch映射注册插件工厂标识为Elasticsearch。其核心职责是把上游 Source/Transform 输出的行数据序列化为 Elasticsearch Bulk API 请求批量提交到目标索引。能力特性对应 connector-v2-featuresCDC变更数据捕获支持可处理 INSERT、UPDATE、DELETE 事件流Exactly-Once精确一次未声明支持。引擎与版本支持官方文档声明支持的 Elasticsearch 版本为 2.x 且 8.x。源码 ElasticsearchVersion.java 中枚举了ES2 / ES5 / ES6 / ES7 / ES8五个版本档位连接器启动时会调用集群接口解析实际版本并映射到对应档位同时兼容 OpenSearch。底层使用elasticsearch-rest-client7.5.1见 pom.xml通过 HTTP REST 协议与集群通信。快速上手最简单配置将Elasticsearch声明为 Sink 只需两个必填参数集群地址hosts与目标索引index。sink { Elasticsearch { hosts [localhost:9200] index seatunnel-${age} } }结合一个真实可运行的作业示例参考仓库 E2E 测试配置 fakesource_to_elasticsearch_multi_sink.conf完整的config文件写法如下env { parallelism 1 job.mode BATCH } source { FakeSource { schema { fields { id int name string age int } } rows [ { kind INSERT, fields [1, Tom, 20] } ] } } transform { } sink { Elasticsearch { hosts [localhost:9200] index seatunnel-${age} schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }配置保存后通过 SeaTunnel 命令行即可提交作业-c指定配置文件-e local使用本地引擎。参数总览下表完整继承自 Elasticsearch.md 官方文档同时依据 SinkConfig.java 与 EsClusterConnectionConfig.java 中的选项定义进行了字段一致性校验名称类型是否必填默认值hostsarray是-indexstring是-schema_save_modestring是CREATE_SCHEMA_WHEN_NOT_EXISTdata_save_modestring是APPEND_DATAindex_typestring否空primary_keyslist否空key_delimiterstring否_usernamestring否空passwordstring否空max_retry_countint否3max_batch_sizeint否10tls_verify_certificateboolean否truetls_verify_hostnameboolean否truetls_keystore_pathstring否-tls_keystore_passwordstring否-tls_truststore_pathstring否-tls_truststore_passwordstring否-common-options-否-必填项的定义同样体现在工厂类的OptionRule中ElasticsearchSinkFactory.java 通过.required(HOSTS, INDEX, SinkConfig.SCHEMA_SAVE_MODE, SinkConfig.DATA_SAVE_MODE)声明了 4 个必填参数其余全部为可选参数。核心参数详解hosts [array]Elasticsearch 集群的 HTTP 地址列表格式为host:port支持指定多个节点以实现连接冗余与负载分担例如[host1:9200, host2:9200]。若集群启用了 HTTPS地址需要写成https://host:9200详见下文 TLS 章节。源码 EsRestClient.java 会将其逐个解析为HttpHost并构建RestClient同时设置连接请求超时10 秒与 Socket 超时5 分钟。index [string]目标索引名支持包含字段名变量格式为seatunnel_${age}这种字段名占位符写法被引用的字段必须存在于 SeaTunnel 行数据中如果字段不存在则该变量不会被替换索引会被当作普通索引名处理。动态索引的底层实现位于 IndexSerializerFactory.java它使用正则\\$\\{(.*?)\\}提取索引名中的所有占位符若存在则创建VariableIndexSerializer否则创建FixedValueIndexSerializer固定索引名。在 VariableIndexSerializer.java 中有三个值得注意的细节变量值取自行中对应字段的toString()若字段值为null则替换为字符串null最终索引名会执行toLowerCase()强制转为小写Elasticsearch 索引名本身也要求小写。sink { Elasticsearch { hosts [localhost:9200] index seatunnel-${age} } }index_type [string]Elasticsearch 索引类型type。在 Elasticsearch 6 及以上版本中建议不要指定因为 6.x 之后 type 概念已逐步废弃7.x 起仅保留_doc。源码 IndexTypeSerializerFactory.java 根据集群实际版本决定行为集群为 OpenSearch直接不写入_type字段ES 2.x / 5.x必须携带 type未配置时自动使用默认值stDEFAULT_TYPE由 RequiredIndexTypeSerializer.java 在 Bulk 元数据中写入_type: typeES 6.x只有显式配置了非空index_type才会写入ES 7.x / 8.x一律不写入。primary_keys [list]用于生成文档_id的主键字段列表这是 CDC 场景的必填选项。配置后序列化器会根据这些字段的值拼接出文档 ID未配置时_id不指定由 Elasticsearch 自动生成。其实现位于 KeyExtractor.java当primary_keys为 null 时 keyExtractor 直接返回null否则按字段顺序取出各字段值并用key_delimiter连接。注意其中 ROW/ARRAY/MAP 类型的字段不允许作为主键DATE/TIME/TIMESTAMP 会按toString()格式输出。key_delimiter [string]复合主键的连接分隔符默认_。例如将分隔符配置为$三个主键KEY1、KEY2、KEY3生成的文档_id为KEY1$KEY2$KEY3。该参数定义于 SinkConfig.java默认值_。username / password [string]X-Pack 安全认证的账号密码。连接器通过BasicCredentialsProvider注入UsernamePasswordCredentials到 HTTP 客户端用于访问开启了安全认证的集群如 Elasticsearch 默认的elastic超级用户。两者均为可选配置了username后通常需同时配置password。max_retry_count [int]单次 Bulk 请求的最大重试次数默认 3。该值直接传递给 ElasticsearchSinkWriter.java 中构造的RetryMaterial重试策略为「始终重试」exception - true每次重试间隔 200msDEFAULT_SLEEP_TIME_MS直至达到最大次数若最终仍失败则抛出ElasticsearchConnectorException并导致作业失败。max_batch_size [int]单个 Bulk 批次的最大文档数默认 10。写入器内部维护一个请求列表当累积的行数达到max_batch_size时立即触发一次批量提交bulkEsWithRetry同时在prepareCommit()checkpoint 提交阶段与close()关闭阶段也会对残留数据执行最终刷写保证数据不丢失。TLS 相关参数参数类型默认值说明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-信任库对应的密钥口令在 EsRestClient.java 的createInstance中有两个值得注意的联动逻辑仅当tls_verify_certificate true时才会读取tls_keystore_path、tls_keystore_password、tls_truststore_path、tls_truststore_password四个参数关闭证书校验后这些参数会被忽略tls_verify_hostname false时使用NoopHostnameVerifier跳过主机名校验tls_verify_certificate false时使用TrustAllStrategy信任所有证书。common options公共参数Sink 插件的公共参数主要包含source_table_name与result_table_name等数据管道串联选项详见 Sink Common Options。当作业中仅有一个 Source、一个 Transform、一个 Sink 时无需指定当任一算子数量大于 1 时必须为每个连接器显式指定source_table_name/result_table_name。schema_save_mode目标表结构预处理策略在同步任务启动之前schema_save_mode决定对目标端已存在的索引表结构采取何种处理方案。可选值RECREATE_SCHEMA表不存在时创建表已存在时先删除再重建CREATE_SCHEMA_WHEN_NOT_EXIST默认表不存在时创建表已存在时跳过ERROR_WHEN_SCHEMA_NOT_EXIST表不存在时直接报错。data_save_mode目标端存量数据处理策略data_save_mode决定在同步任务启动之前对目标端已存在的数据采取何种处理方案。可选值DROP_DATA保留库表结构删除存量数据APPEND_DATA默认保留库表结构保留存量数据新数据追加写入ERROR_WHEN_DATA_EXISTS目标端已有数据时报错。从源码看SinkConfig.java 中DATA_SAVE_MODE的可选值被限定为上述三种singleChoiceSCHEMA_SAVE_MODE则取SchemaSaveMode枚举。在 ElasticsearchSink.java 中Sink 实现了SupportSaveMode接口通过getSaveModeHandler()动态发现同名的CatalogFactory创建ElasticSearchCatalog再用DefaultSaveModeHandler将schema_save_mode与data_save_mode翻译为实际的建表 / 删表 / 清数动作。典型配置sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }写入实现原理Bulk 批量提交与重试ElasticsearchSink在 ElasticsearchSink.java 中创建 ElasticsearchSinkWriter.java其写入流程为逐行序列化每收到一行SeaTunnelRow先用ElasticsearchRowSerializer序列化为 Bulk 请求文本见下一节放入内存列表批量触发当列表大小达到max_batch_size默认 10时调用bulkEsWithRetry一次性提交带重试提交RetryUtils.retryWithException包裹提交逻辑将列表用\n连接成请求体发给esRestClient.bulk(...)若响应中errorstrue部分文档失败抛出异常触发重试最多max_retry_count次成功后清空列表生命周期兜底prepareCommit()checkpoint 提交时和close()writer 关闭时都会再执行一次bulkEsWithRetry确保缓冲区残留数据全部落库随后关闭EsRestClient。这里有一个明显的性能权衡点默认max_batch_size 10意味着每积累 10 行就发起一次 HTTP 请求对于高吞吐场景建议结合实际数据量调大该值例如 1000 ~ 5000以显著降低请求次数、提升写入吞吐同时配合max_retry_count控制失败容忍度。CDC 事件语义INSERT / UPDATE / DELETE 的序列化与 _id 生成连接器支持 CDC 事件流核心逻辑在 ElasticsearchRowSerializer.java。序列化器按行的RowKind分派处理INSERT/UPDATE_AFTER→ 执行upsert若存在主键_id生成{ update: {_index: ..., _id: ...} }{ doc: {文档}, doc_as_upsert: true }两条 NDJSON 行实现「存在即更新、不存在即插入」若无主键则生成{ index: {...} } 文档 JSON走普通索引写入UPDATE_BEFORE/DELETE→ 执行delete生成{ delete: {_index: ..., _id: ...} }按主键删除文档其他 RowKind → 抛出UNSUPPORTED_OPERATION异常。同时在 ElasticsearchSinkWriter.java 的write方法中UPDATE_BEFORE行会被直接跳过因为UPDATE_AFTER的 upsert 已经覆盖了更新语义。文档 JSON 的构建逻辑toDocumentMap会递归展开嵌套的SeaTunnelRow结构化类型并针对 JDK 8 时间类型Temporal执行toString()转换Jackson 默认不支持直接序列化这些类型Map/List 内嵌套的值也会递归转换。CDC 场景的必选配置示例sink { Elasticsearch { hosts [localhost:9200] index seatunnel-${age} # cdc required options primary_keys [key1, key2, ...] } }SSL/TLSHTTPS 安全连接配置当集群启用 HTTPS 时可按需组合以下配置。文档给出了四种典型场景关闭证书校验仅跳过证书链校验适合自签名证书快速联调sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_verify_certificate false } }关闭主机名校验证书合法但与主机名不匹配时使用sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_verify_hostname false } }启用证书校验推荐生产用法加载本地密钥库sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_keystore_path ${your elasticsearch home}/config/certs/http.p12 tls_keystore_password ${your password} } }需要说明的是生产环境应保持tls_verify_certificate与tls_verify_hostname的默认值true并通过tls_keystore_path/tls_truststore_path配置受信证书而非直接关闭校验。多表写入与测试验证从 ElasticsearchSink.java 可以看到Sink 还实现了SupportMultiTableSink接口支持多表场景。E2E 测试配置 fakesource_to_elasticsearch_multi_sink.conf 演示了 FakeSource 同时产出st_index5、st_index6两张表、由同一个 Elasticsearch Sink 写入的场景其中索引名使用了index ${table_name}动态变量对应多表框架注入的table_name字段sink { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index ${table_name} index_type st schema_save_modeCREATE_SCHEMA_WHEN_NOT_EXIST data_save_modeAPPEND_DATA } }此外elasticsearch_source_and_sink.conf 与 elasticsearch_source_without_schema_and_sink.conf 分别覆盖了「带 Schema 的 ES → ES 全链路」与「不带 Schema 的 ES → ES」两类场景测试断言ElasticsearchIT.java通过查询目标索引并比对文档内容验证了写入结果的正确性。生产实践要点与限制索引名强制小写动态索引最终会toLowerCase()请确保目标索引名符合小写规范字段缺失时的索引名若动态变量引用的字段不在行数据中占位符不会被替换索引名保持原样容易被误认为普通索引ES 6 及以上不要配置index_type7.x / 8.x 已不再支持自定义 type配置后也不会写入由 IndexTypeSerializerFactory.java 自动忽略主键类型限制primary_keys指定的字段不能是 ROW / ARRAY / MAP 类型日期与时间类型会以toString()形式参与_id拼接批量大小权衡默认max_batch_size 10偏保守高吞吐场景应调大max_retry_count 3与 200ms 固定重试间隔可覆盖大多数瞬时故障版本适用范围支持 Elasticsearch 2.x ~ 8.x并兼容 OpenSearch更早或更新的版本不在官方支持范围内。变更记录2.2.0-beta2022-09-26新增 Elasticsearch Sink 连接器next version支持 CDC 写入 DELETE / UPDATE / INSERT 事件PR #3673支持 HTTPS 协议并兼容 OpenSearchPR #3997。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Hudi Sink 连接器详解配置参数、多表写入与 CDC 实战指南SeaTunnel Hudi Sink 连接器详解配置参数、多表写入与 CDC 实战指南 本指南以 SeaTunnel 仓库中 Hudi Sink 官方文档数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Kudu Sink 连接器实战指南参数详解、CDC 写入与多表路由SeaTunnel Kudu Sink 连接器实战指南参数详解、CDC 写入与多表路由 本指南以 SeaTunnel 仓库中 Kudu Sink 官方文档 h数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Elasticsearch Sink 连接器实战指南Bulk 写入、CDC 语义、认证与 TLS 配置详解SeaTunnel Elasticsearch Sink 连接器实战指南Bulk 写入、CDC 语义、认证与 TLS 配置详解 本文系统讲解 SeaTunne数据集成ETL大数据批处理流处理变更数据捕获上一篇Windows HEIC缩略图插件三步注册让资源管理器直接预览iPhone照片下一篇CSDN博客下载器使用教程免费批量离线保存技术文章的完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表