ARTICLE DETAIL

资讯详情

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

SeaTunnel Elasticsearch Sink 连接器实战指南:配置、认证、CDC 与多表写入

SeaTunnel Elasticsearch Sink 连接器实战指南:配置、认证、CDC 与多表写入 SeaTunnel Elasticsearch Sink 连接器实战指南配置、认证、CDC 与多表写入【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇指南以 Apache SeaTunnelseatunnel-connectors-v2/connector-elasticsearch中的 Elasticsearch Sink 连接器为核心系统讲解如何将数据写入 Elasticsearch 集群。文章涵盖连接器的全部核心选项hosts、index、schema/data save mode、主键、批量与重试、TLS、向量字段等、三种认证方式basic、api_key、api_key_encoded、基于 Bulk API 的写入实现原理以及 CDC 变更数据捕获、多表写入、定时刷新、模式演变等实战场景帮助读者从配置到原理完整掌握该 Sink 插件的使用与调优。描述与主要特性SeaTunnel 的 Elasticsearch Sink 插件用于将上游数据输出到Elasticsearch集群支持 Elasticsearch 2.x 到 8.x 的版本范围该范围同时记录于 ElasticsearchRowSerializer 的类注释中。依据文档 connector-v2-features 的能力矩阵该连接器支持的特性如下变更数据捕获CDC支持可将 CDC 事件中的 INSERT / UPDATE / DELETE 正确映射为 Elasticsearch 文档操作多表写入支持可在单个 Sink 配置下将多张表写入各自索引定时刷新支持由 Zeta 引擎注入定时刷新信号精确一次Exactly-Once不支持。其中定时刷新仅由 Zeta 引擎支持Spark 与 Flink 运行时不注入FlushSignal。选项总览连接器选项定义集中在 ElasticsearchSinkOptions 与 ElasticsearchBaseOptions 两个配置类中。完整选项表如下名称类型是否必须默认值hostsarray是-indexstring是-schema_save_modestring否CREATE_SCHEMA_WHEN_NOT_EXISTdata_save_modestring否APPEND_DATAindex_typestring否空primary_keyslist否空key_delimiterstring否_auth_typestring否basicusernamestring否空passwordstring否空auth.api_key_idstring否-auth.api_keystring否-auth.api_key_encodedstring否-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-否-vectorization_fieldsarray否-vector_dimensionsint否0multi_table_sink_replicaint否1在 ElasticsearchSinkFactory 的optionRule()中可以看到这些选项的校验规则hosts、index、schema_save_mode、data_save_mode为必填项其余为可选项且认证相关选项之间存在条件约束详见下文认证章节。hosts [array]Elasticsearch 集群的 HTTP 地址格式为host:port允许指定多个主机以实现连接层面的负载均衡与故障切换例如[host1:9200, host2:9200]。生产环境建议配置集群中多个节点的地址。index [string]Elasticsearch 的索引名称。索引名支持字段名变量例如seatunnel_${age}此时要求对应的字段必须出现在 SeaTunnel Row 中且建议配合schema_save_mode IGNORE使用如果字段不存在则将其视为普通索引名。索引名最终会被转换为小写。从源码实现看IndexSerializerFactory 会根据索引名是否包含${来选择序列化器固定索引由 FixedValueIndexSerializer 直接返回配置的索引名变量索引由 VariableIndexSerializer 在每条记录写入时用当前行中对应字段的值替换${fieldName}占位符字段值为 null 时替换为字符串null。借助该机制可以按字段值动态路由数据到不同的索引例如按年龄、日期分索引。index_type [string]Elasticsearch 索引类型。Elasticsearch 6 及以上版本已废弃多类型概念建议不要指定默认留空。源码中 IndexTypeSerializerFactory 会根据集群信息与类型决定使用 RequiredIndexTypeSerializer需要填充_type还是 NotIndexTypeSerializer。primary_keys [list]主键字段列表用于生成文档_id。这是 CDC 场景的必需选项因为只有具备确定性的_idUPDATE / DELETE 事件才能定位到目标文档。key_delimiter [string]复合主键的分隔符默认_。例如使用$作为分隔符时文档_id将呈现为KEY1$KEY2$KEY3格式。该逻辑由 KeyExtractor 实现它按primary_keys在 Row 中定位字段索引并用分隔符拼接各字段的格式化值。multi_table_sink_replica [int]多表写入时每张表对应的 Sink Writer 副本数。通常保持默认值 1 即可只有单表写入压力较大、需要更高写入并行度时才调大。该选项属于 Sink 连接器通用选项定义于 SinkConnectorCommonOptions。max_retry_count [int]单个 Bulk 请求的最大重试次数默认 3。对应 ElasticsearchSinkWriter 中的RetryMaterial重试间隔固定为 200ms重试条件为任意异常exception - true。max_batch_size [int]单个批次内缓存的最大文档数默认 10。当待写入文档数达到该阈值时会触发一次 Bulk 写入requestEsList.size() maxBatchSize即调用bulkEsWithRetry。在实际生产环境中建议结合吞吐量将该值调大例如数万级别以获得更好的批量写入性能。vectorization_fields [array]需要向量转换的字段名列表字段值应为ByteBuffer类型的嵌入向量Elasticsearch 7.3 及以后版本支持。从 ElasticsearchRowSerializer.convertValue 的实现可以看到ByteBuffer会被转换为Float[]数组若该字段命中vectorization_fields且配置了vector_dimensions则按配置维度截取否则按 buffer 剩余字节数除以 4 计算维度。vector_dimensions [int]向量字段的维度即向量中浮点数个数默认 0。同样要求 Elasticsearch 7.3 及以后版本。common optionsSink 插件通用参数详见 Sink 常用选项。认证方式Elasticsearch 连接器支持三种认证方式配置枚举定义于 AuthTypeEnum具体认证实现位于 client/auth 目录下包括BasicAuthProvider、ApiKeyAuthProvider、ApiKeyEncodedAuthProvider由 AuthenticationProviderFactory 按auth_type创建。auth_type [enum]指定使用的认证方式支持的值basic默认使用用户名和密码的 HTTP 基本认证api_key使用独立的 ID 和密钥的 Elasticsearch API Key 认证api_key_encoded使用编码密钥的 Elasticsearch API Key 认证。未指定时默认使用basic以保持向后兼容。三种枚举值在源码 AuthTypeEnum 中与字符串值一一对应。基本认证basic基本认证使用 HTTP 基本认证通过用户名和密码凭据认证对应 x-pack 的 username / password。username基本认证用户名x-pack 用户名password基本认证密码x-pack 密码。sink { Elasticsearch { hosts [https://localhost:9200] auth_type basic username elastic password your_password index my_index } }API Key 认证api_keyAPI Key 认证提供了比用户名密码更安全的方式使用 Elasticsearch 生成的 API 密钥进行认证auth.api_key_idElasticsearch 生成的 API 密钥 IDauth.api_keyElasticsearch 生成的 API 密钥secret。sink { 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 } }编码 API Key 认证api_key_encodedauth.api_key_encodedBase64 编码的 API 密钥格式为base64(id:api_key)是分别指定auth.api_key_idauth.api_key的替代方式。注意可以使用auth.api_key_idauth.api_key或auth.api_key_encoded但不能同时使用两者。该约束在 ElasticsearchSinkFactory.optionRule 中通过条件规则实现basic要求捆绑username/password且二者非空api_key要求auth.api_key_id与auth.api_key同时非空api_key_encoded要求auth.api_key_encoded非空且满足ApiKeyEncodedFormatValidator格式校验。sink { Elasticsearch { hosts [https://localhost:9200] auth_type api_key_encoded auth.api_key_encoded eW91cl9hcGlfa2V5X2lkOnlvdXJfYXBpX2tleV9zZWNyZXQ index my_index } }写入流程与批量语义Elasticsearch Sink 的写入链路为ElasticsearchSink→ElasticsearchSinkWriter→EsRestClient。行序列化每条 SeaTunnel Row 由 ElasticsearchRowSerializer.serializeRow 根据RowKind转换成对应的 Bulk API 请求行INSERT / UPDATE_AFTER生成index或update操作。若配置了primary_keys即 key 非 null生成{ update : {_index: ..., _id: ...} }\n{ doc : {...}, doc_as_upsert : true }实现 upsert存在则更新、不存在则插入语义否则生成{ index : {_index: ...} }\n{...doc...}追加写入。UPDATE_BEFORE / DELETE生成{ delete : {_index: ..., _id: ...} }删除操作。批量与重试ElasticsearchSinkWriter.write将序列化后的请求行缓存到requestEsList当达到max_batch_size时调用bulkEsWithRetry将请求拼接为 NDJSONString.join(\n, ...)一次性提交给 EsRestClient.bulk。若 Bulk 响应包含错误bulkResponse.isErrors()会抛出异常并按照max_retry_count进行重试。触发刷新的时机从 ElasticsearchSinkWriter 源码可归纳出四个刷新触发点缓存文档数达到max_batch_sizeprepareCommit()checkpoint 提交前Zeta 引擎的定时刷新信号timerFlush注册于context.registerFlushActionWriterclose()时最后一次 flush。语义说明Elasticsearch Sink 当前提供至少一次at-least-once语义。定时刷新不提供基于 2PC 的精确一次语义如果文档 ID 不是确定性的例如未配置primary_keys失败重试可能产生重复写入。因此对数据一致性要求较高的场景如 CDC务必配置确定性的primary_keys以便依赖_id幂等去重。Zeta 定时刷新该引擎级能力仅由 Zeta 支持Spark 和 Flink 不会注入FlushSignal。在 Zeta 中可以在env块配置sink.flush.interval使未达到max_batch_size的待处理 Bulk 请求也能定时写出从而降低数据延迟。该配置项在引擎侧由 ServerConfigOptions 定义并由引擎的任务流生命周期如 SinkFlowLifeCycle向 Sink Writer 注入刷新信号相关行为还有对应的引擎测试用例如 SeaTunnelSourceCollectorFlushSignalTest。env { job.mode STREAMING checkpoint.interval 300000 sink.flush.interval 5000 } sink { Elasticsearch { hosts [localhost:9200] index seatunnel-index max_batch_size 10000 } }上述配置表示流式任务每 5 秒定时刷新一次待处理的 Bulk 请求即使文档数未达到max_batch_size10000也会写出。表结构与数据处理策略schema_save_mode在启动同步任务之前针对目标侧已有的表结构索引结构选择不同的处理方案RECREATE_SCHEMA当索引不存在时会创建当索引已存在时会删除并重建CREATE_SCHEMA_WHEN_NOT_EXIST默认当索引不存在时会创建当索引已存在时跳过创建ERROR_WHEN_SCHEMA_NOT_EXIST当索引不存在时将抛出错误IGNORE忽略对索引的处理。data_save_mode在启动同步任务之前针对目标侧已存在的数据选择不同的处理方案DROP_DATA保留索引结构删除数据APPEND_DATA默认保留索引结构保留数据追加写入ERROR_WHEN_DATA_EXISTS当索引中已有数据时抛出错误。这两个模式由 ElasticsearchSink.getSaveModeHandler 实现它会通过CatalogFactory机制发现 ElasticSearchCatalog构造DefaultSaveModeHandler在任务启动前执行建表/删表/校验等预处理。注意data_save_mode在 ElasticsearchSinkOptions 中通过singleChoice限制了可选值为DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS三种。典型使用示例简单示例sink { Elasticsearch { hosts [localhost:9200] index seatunnel-${age} schema_save_modeIGNORE } }该示例演示了动态索引每条记录根据age字段的值路由到seatunnel-18、seatunnel-30等不同索引。多表写入sink { Elasticsearch { hosts [localhost:9200] index ${table_name} schema_save_modeIGNORE multi_table_sink_replica 1 } }通过${table_name}变量将多张源表的数据写入各自同名的索引。从源码看多表场景下多个 Writer 通过 ElasticsearchMultiTableResourceManager 共享同一个EsRestClient连接组Writer 在setMultiTableResourceManager中注入共享客户端避免每张表都建立独立连接。向量转换vector datasink { Elasticsearch { hosts [localhost:9200] index ${table_name} schema_save_modeIGNORE vectorization_fields [review_embedding] vector_dimensions 1024 } }将review_embedding字段中的ByteBuffer向量转换为 1024 维的Float[]写入 Elasticsearch需要 ES 7.3且目标索引应预先定义好dense_vector映射。变更数据捕获CDC事件sink { Elasticsearch { hosts [localhost:9200] index seatunnel-${age} schema_save_modeIGNORE # CDC required options primary_keys [key1, key2, ...] } }CDC 场景必须配置primary_keys用于生成确定性的文档_id从而支持 UPDATE / DELETE 事件的准确定位。CDC 事件多表写入sink { Elasticsearch { hosts [localhost:9200] index ${table_name} schema_save_modeIGNORE primary_keys [${primary_key}] } }注意这里的primary_keys支持${primary_key}形式的动态占位配合多表 CDC 时按源表各自的主键字段生成_id。TLS/SSL 配置当 Elasticsearch 集群启用 HTTPS 时需要配置 TLS 相关选项。相关参数在 ElasticsearchBaseOptions 中定义实际 SSL 上下文构建逻辑位于 SSLUtils。tls_verify_certificate为 HTTPS 端点启用证书验证默认 truetls_verify_hostname为 HTTPS 端点启用主机名验证默认 truetls_keystore_path指向 PEM 或 JKS 密钥存储的路径运行 SeaTunnel 的操作系统用户必须能够读取该文件tls_keystore_password指定的密钥存储的密钥密码tls_truststore_path指向 PEM 或 JKS 信任存储的路径运行 SeaTunnel 的操作系统用户必须能够读取该文件tls_truststore_password指定的信任存储的密钥密码。SSL 禁用证书验证sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_verify_certificate false } }SSL 禁用主机名验证sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch tls_verify_hostname false } }SSL 启用证书验证通过设置tls_keystore_path与tls_keystore_password指定客户端证书路径及密码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} } }配置表生成策略通过将schema_save_mode配置为CREATE_SCHEMA_WHEN_NOT_EXIST支持索引不存在时自动创建sink { Elasticsearch { hosts [https://localhost:9200] username elastic password elasticsearch schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }模式演变Schema EvolutionCDC 采集支持有限数量的模式更改目前支持的模式更改包括添加列。从源码看ElasticsearchSink.supports 返回SchemaChangeType.ADD_COLUMN而 ElasticsearchSinkWriter.applySchemaChange 在收到AlterTableAddColumnEvent时会通过ElasticSearchTypeConverter将 SeaTunnel 列类型转换为 ES 类型并调用esRestClient.addField(index, ...)动态为索引添加字段。模式演变完整示例MySQL-CDC 动态加列同步到 Elasticsearchenv { # You can set engine configuration here parallelism 5 job.mode STREAMING checkpoint.interval 5000 read_limit.bytes_per_second7000000 read_limit.rows_per_second400 } source { MySQL-CDC { server-id 5652-5657 username st_user_source password mysqlpw table-names [shop.products] url jdbc:mysql://mysql_cdc_e2e:3306/shop schema-changes.enabled true } } sink { Elasticsearch { hosts [https://elasticsearch:9200] username elastic password elasticsearch tls_verify_certificate false tls_verify_hostname false index schema_change_index index_type _doc schema_save_modeCREATE_SCHEMA_WHEN_NOT_EXIST data_save_modeAPPEND_DATA } }前提条件MySQL-CDC Source 需开启schema-changes.enabled true才能在源端 DDL添加列发生时将AlterTableAddColumnEvent事件传递给 Elasticsearch Sink 完成索引字段的同步添加。测试与验证仓库中提供了针对该连接器的端到端测试可作为配置与行为的参考验证ElasticsearchIT基础读写集成测试ElasticsearchAuthIT认证相关集成测试ElasticsearchSchemaChangeIT模式演变加列集成测试ElasticsearchSinkWriterTestWriter 单元测试。变更日志该连接器的完整版本变更记录见 connector-elasticsearch 变更日志。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表