- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
SeaTunnel 内置的ElasticsearchSink 插件用于将 SeaTunnel 管道中的行数据(SeaTunnelRow)批量写入 Elasticsearch 集群,支持动态索引名、主键文档_id生成、CDC(Change 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-features):
- CDC(变更数据捕获):支持,可处理 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 中的选项定义进行了字段一致性校验:
| 名称 | 类型 | 是否必填 | 默认值 |
|---|---|---|---|
| hosts | array | 是 | - |
| index | string | 是 | - |
| schema_save_mode | string | 是 | CREATE_SCHEMA_WHEN_NOT_EXIST |
| data_save_mode | string | 是 | APPEND_DATA |
| index_type | string | 否 | (空) |
| primary_keys | list | 否 | (空) |
| key_delimiter | string | 否 | _ |
| username | string | 否 | (空) |
| password | string | 否 | (空) |
| max_retry_count | int | 否 | 3 |
| max_batch_size | int | 否 | 10 |
| tls_verify_certificate | boolean | 否 | true |
| tls_verify_hostname | boolean | 否 | true |
| tls_keystore_path | string | 否 | - |
| tls_keystore_password | string | 否 | - |
| tls_truststore_path | string | 否 | - |
| tls_truststore_password | string | 否 | - |
| 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,未配置时自动使用默认值
st(DEFAULT_TYPE),由 RequiredIndexTypeSerializer.java 在 Bulk 元数据中写入"_type": type; - ES 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),每次重试间隔 200ms(DEFAULT_SLEEP_TIME_MS),直至达到最大次数;若最终仍失败,则抛出ElasticsearchConnectorException并导致作业失败。
max_batch_size [int]
单个 Bulk 批次的最大文档数,默认 10。写入器内部维护一个请求列表,当累积的行数达到max_batch_size时立即触发一次批量提交(bulkEsWithRetry);同时在prepareCommit()(checkpoint 提交阶段)与close()(关闭阶段)也会对残留数据执行最终刷写,保证数据不丢失。
TLS 相关参数
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| tls_verify_certificate | boolean | true | 是否校验 HTTPS 端点的证书链 |
| tls_verify_hostname | boolean | true | 是否校验 HTTPS 端点的主机名 |
| tls_keystore_path | string | - | PEM 或 JKS 格式的密钥库路径,运行 SeaTunnel 的操作系统用户必须可读 |
| tls_keystore_password | string | - | 密钥库对应的密钥口令 |
| tls_truststore_path | string | - | PEM 或 JKS 格式的信任库路径,运行 SeaTunnel 的操作系统用户必须可读 |
| tls_truststore_password | string | - | 信任库对应的密钥口令 |
在 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的可选值被限定为上述三种(singleChoice),SCHEMA_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(...);若响应中errors=true(部分文档失败),抛出异常触发重试,最多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/TLS:HTTPS 安全连接配置
当集群启用 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_mode"="CREATE_SCHEMA_WHEN_NOT_EXIST" "data_save_mode"="APPEND_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_type:7.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-beta(2022-09-26):新增 Elasticsearch Sink 连接器;
- next version:
- 支持 CDC 写入 DELETE / UPDATE / INSERT 事件(PR #3673);
- 支持 HTTPS 协议并兼容 OpenSearch(PR #3997)。
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
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大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考