尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SeaTunnel 数据 Sink 连接器选型指南:写入模式、Save Mode 与常用参数实战

SeaTunnel 数据 Sink 连接器选型指南:写入模式、Save Mode 与常用参数实战 SeaTunnel 数据 Sink 连接器选型指南写入模式、Save Mode 与常用参数实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel当你搭建 SeaTunnel 数据集成任务时第一个要回答的问题往往是“数据最终要写到哪里”。本篇基于仓库中的 Sink 连接器总览、Sink 常用选项 与 Sink 写入模式与 Save Mode 整理先给出选择 Sink 前的四项确认清单再深入讲解plugin_input、parallelism等常用参数、generate_sink_sql与query两种互斥写入模式、schema_save_mode/data_save_mode的语义边界并结合seatunnel-api与 JDBC 连接器的源码实现帮助你在配置任何具体 Sink 之前建立完整的选型决策框架。一、选择 Sink 前先确认四件事SeaTunnel 仓库共收录了近百个 Sink 连接器页面docs/zh/connectors/sink/ 目录下共 98 个文档目标系统覆盖关系数据库、消息队列、对象存储、搜索引擎、NoSQL 与各类 SaaS 平台。面对如此多的选择官方文档给出的选型顺序是先匹配目标系统再依次确认写入保证、表结构行为、驱动依赖以及连接器自身的投递约束。具体展开为四项检查清单最终目标系统以及表、对象或主题的落点形态数据要落到哪张表、哪个主题、哪个目录或对象路径目标端的数据组织形式直接决定选型例如文件 Sink 与数据库 Sink 的参数体系完全不同。任务需要至少一次at-least-once、精确一次exactly-once还是幂等写入不同 Sink 对精确一次的实现方式不同例如 JDBC Sink 提供基于事务的JdbcExactlyOnceSinkWriter见 seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcExactlyOnceSinkWriter.java而 Console 这类只输出日志的 Sink 则不提供任何写入语义保证。是否依赖额外驱动、SDK 或云鉴权配置JDBC 系列 Sink 需要数据库驱动对象存储类 Sink 需要 AccessKey/SecretKey 或云厂商凭证依赖的加载方式受 连接器依赖隔离加载机制 约束。该连接器是否支持你要使用的写入模式例如schema_save_mode、data_save_mode并非所有 Sink 都支持最终要以你所用版本对应连接器页面的参数表为准仓库中每个连接器页面均以“主要特性”勾选框标注其对精确一次、CDC、多表写入等特性的支持情况。二、Sink 连接器在仓库中的组织方式从仓库结构可以确认 Sink 连接器的三类证据文档层docs/zh/connectors/sink/ 下每个连接器一个页面统一包含“支持连接器版本、支持的引擎、描述、主要特性、核心参数、示例配置”等章节实现层每个连接器模块如 seatunnel-connectors-v2/connector-jdbc/实现seatunnel-api中定义的 Sink 接口注册层plugin-mapping.properties 中通过seatunnel.sink.PluginName connector-module的形式声明 Sink 插件到依赖目录的映射当前共注册了 84 个seatunnel.sink.*条目例如seatunnel.sink.Console connector-console seatunnel.sink.Kafka connector-kafka seatunnel.sink.Jdbc connector-jdbc这个映射值同时决定了依赖隔离加载时的插件目录名${SEATUNNEL_HOME}/plugins/connector-module/详见 连接器依赖隔离加载机制。从源码结构看所有 Sink 的通用能力都收敛在 seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/ 包中接口/类职责SeaTunnelSinkSink 连接器顶层接口SinkWriter/SupportMultiTableSinkWriter单表/多表写入器SinkCommitter/SinkAggregatedCommitter两阶段提交checkpoint 提交组件是精确一次语义的关键SupportSaveModeSaveModeHandler/DefaultSaveModeHandlerSave Mode 处理 SPI 与默认实现SchemaSaveMode/DataSaveMode两类 Save Mode 的枚举定义其中 SchemaSaveMode.java 与 DataSaveMode.java 两个枚举的取值与文档描述完全一致可作为配置合法值的权威依据。三、Sink 常用选项详解所有 Sink 连接器共享一组通用参数定义见 Sink 常用选项名称类型是否需要默认值说明plugin_inputstring否-指定当前 Sink 处理哪个上游数据集parallelismint否-未指定时回落到env中的parallelism指定后覆盖全局值metadata_datasource_idstring否-指定后连接器从外部元数据服务获取连接信息URL、用户名、密码替代直接配置迁移提示旧配置名source_table_name已过时请使用新名称plugin_input。数据集路由规则plugin_input不指定plugin_input时Sink 默认消费配置文件中上一个插件输出的数据集dataset指定plugin_input时Sink 消费该参数对应名称的数据集。完整可运行示例以下示例展示了数据集路由的典型用法一个FakeSourceStream输出数据集fake经两个Filter拆分为fake_name与fake_age两个数据集最后由两个ConsoleSink 分别消费source { FakeSourceStream { parallelism 2 plugin_output fake field_name name,age } } transform { Filter { plugin_input fake fields [name] plugin_output fake_name } Filter { plugin_input fake fields [age] plugin_output fake_age } } sink { Console { plugin_input fake_name } Console { plugin_input fake_age } }路由规则的简化结论如果作业只有一个 source、零个或多个 transform、一个 sink则不需要为连接器指定plugin_input/plugin_output如果 source、transform、sink 中任意一侧的算子数量大于 1则必须为作业中每个连接器显式指定plugin_input和plugin_output否则数据集会相互串流。四、写入模式与 Save Mode两组最易混淆的决策配置 Sink 时有两个维度最容易混淆详见 Sink 写入模式与 Save Mode写入模式决定 SeaTunnel如何把每一行数据写到目标端Save Mode决定 SeaTunnel 在开始写入数据前如何处理目标端已存在的表、索引、目录或数据。4.1 快速决策表目标优先选择说明让 SeaTunnel 为 JDBC 目标端生成 INSERT / UPSERT / UPDATE / DELETE SQLgenerate_sink_sql true并配置database、table通常还要配置primary_keysJDBC Sink 能力SeaTunnel 能解析目标 Catalog 表时也可执行 save mode 和自动建表完全控制 JDBC 写入 SQLquery INSERT ... VALUES (?, ...)不要与generate_sink_sql true同时配置此模式下不执行schema_save_mode、data_save_mode或custom_sql目标表不存在时自动创建或不存在时报错schema_save_mode仅适用于显式暴露 save mode 参数、且能通过 Catalog 创建或检查目标端的 Sink写入前保留、清空或检查目标端已有数据data_save_mode支持值取决于具体连接器File Sink 通常只支持DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS写入数据前先执行一条自定义 SQLdata_save_mode CUSTOM_PROCESSING和custom_sql仅适用于同时暴露这两个参数的连接器这是写入前钩子不是逐行写入 SQLJDBC Sink 使用数据库原生 Upsertgenerate_sink_sql true、primary_keys、enable_upsert true没有可用主键或唯一键时自动生成的 SQL 会退化为普通 INSERT写入对象存储或文件系统查看具体 File Sink 参数表File Sink 不使用generate_sink_sql4.2 JDBC Sinkgenerate_sink_sql与query互斥JDBC Sink 提供两种互斥写入模式模式必需参数是否执行 Save Mode典型场景自动生成 SQLgenerate_sink_sql true、database通常还有table是前提是能解析目标 Catalog 表大多数数据库写入、CDC 写入、自动建表、upsert/update/delete自定义 SQLquery INSERT ... VALUES (?, ...)否需要完全控制目标 SQL并接受跳过 save mode 处理使用generate_sink_sql true且需要处理 UPDATE、DELETE 或 upsert 记录时请配置primary_keys。若未显式配置SeaTunnel 会按以下顺序尝试继承上游 Catalog 元数据中的主键 → 取第一组唯一键 → 都拿不到则退化为普通 INSERT。这一点也可以从 JdbcSink.java 所在的 JDBC Sink 实现目录结构中印证——该目录下同时存在JdbcSinkWriter逐行写入与JdbcExactlyOnceSinkWriter事务保证两个写入器以及JdbcSinkCommitter/JdbcSinkAggregatedCommitter两个提交组件分别对应至少一次与精确一次两种投递保证。4.3 Save Mode 语义与源码枚举一一对应schema_save_mode控制写入前如何处理目标结构值行为RECREATE_SCHEMA目标不存在时创建已存在时删除后重建CREATE_SCHEMA_WHEN_NOT_EXIST仅在目标不存在时创建ERROR_WHEN_SCHEMA_NOT_EXIST目标不存在时报错IGNORE跳过结构处理data_save_mode控制写入前如何处理目标端已有数据值行为DROP_DATA保留结构并清空已有数据APPEND_DATA保留已有数据并追加写入CUSTOM_PROCESSING写入前执行custom_sql仅适用于同时暴露这两个参数的连接器ERROR_WHEN_DATA_EXISTS发现已有数据时报错以上 8 个枚举值与 SchemaSaveMode.java、DataSaveMode.java 中的定义完全一致。对于数据库 Sink处理对象通常是表对于文件 Sink则是路径或目录。4.4 各系列 Sink 的支持边界JDBC 系列 SinkJDBC 及 MySQL、PostgreSQL、Oracle、SQL Server 等使用同一套写入模式支持generate_sink_sql与query自动生成 SQL 模式下支持schema_save_mode和data_save_modecustom_sql只有在 save mode 处理真正执行时才会执行enable_upsert只有在 SeaTunnel 拿到可用主键或唯一键后才有意义。完整参数见 JDBC Sink。Doris Sink支持schema_save_mode、data_save_mode、custom_sql和save_mode_create_template但不使用 JDBC 的generate_sink_sql。若要处理 CDC DELETE 事件还需 Doris 侧支持删除能力并按场景配置sink.enable-delete详见 Doris Sink。File 与对象存储 Sink写的是文件因此不使用generate_sink_sql、query或数据库 upsert。各文件连接器暴露 save mode 参数的情况如下Connector是否暴露 Save Mode 参数说明LocalFile是处理本地目录和文件HdfsFile是处理 HDFS 目录和文件FtpFile是处理 FTP 目录和文件SftpFile是处理 SFTP 目录和文件S3File是通过 File Sink save mode 流程处理 S3 路径和对象OssFile是通过 File Sink save mode 流程处理 OSS 路径和对象ObsFile否当前 sink option rule 未暴露schema_save_mode/data_save_modeCosFile否同上BosFile否同上如果某个文件连接器页面没有列出schema_save_mode或data_save_mode不要默认认为它可以接收这些参数。4.5 三个典型配置示例JDBC 自动生成 SQL 并使用 Save Modesink { Jdbc { url jdbc:postgresql://localhost:5432/sales driver org.postgresql.Driver username postgres password change_me generate_sink_sql true database sales table public.orders primary_keys [id] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }JDBC 自定义 SQL不执行 Save Modesink { Jdbc { url jdbc:mysql://localhost:3306/sales driver com.mysql.cj.jdbc.Driver username root password change_me query INSERT INTO orders(id, amount) VALUES (?, ?) } }该模式下 JDBC Sink 只通过query逐行写入不会执行schema_save_mode、data_save_mode或custom_sql。S3File 写入前清空已有数据sink { S3File { path /warehouse/orders bucket s3a://example-bucket fs.s3a.endpoint s3.amazonaws.com fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider access_key ... secret_key ... file_format_type json schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode DROP_DATA } }4.6 故障排查配了generate_sink_sql true但仍然只是 INSERT检查 SeaTunnel 是否拿到了可用 key。需要 upsert、update 或 delete 时建议显式配置primary_keys。JDBC Sink 的custom_sql没有执行检查是否配置了query。JDBC 自定义 query 模式不会执行 save mode 处理因此会跳过custom_sql。File Sink 不接受data_save_mode核对该连接器参数表。S3File、OssFile、HdfsFile、FtpFile、SftpFile、LocalFile暴露文件 save mode 参数ObsFile、CosFile、BosFile当前未暴露。只想建表、不想抽取数据Save mode 是 Sink 作业写入前的一部分SeaTunnel 目前没有通过schema_save_mode提供独立的“只执行 DDL”模式。即使作业没有数据Sink 仍可能完成初始化但这不能替代专门的 schema 管理流程。五、依赖隔离Sink 驱动与 SDK 的正确放置方式上文第 3 项检查驱动/SDK/鉴权依赖与加载机制强相关。SeaTunnel 为每个连接器提供依赖隔离加载连接器自身的依赖 jar 需放置在${SEATUNNEL_HOME}/plugins/connector-xxx/目录下子目录名取自 plugin-mapping.properties 中的 valueSeaTunnel 启动时只加载对应目录的 jar避免不同连接器依赖冲突。从源码结构看该机制由 seatunnel-plugin-discovery/ 模块中的插件发现逻辑实现任务日志中AbstractPluginDiscovery打印的find connector jar and dependency for PluginIdentifier{...}行可用于验证每个 Sink 只加载了自己的依赖。需要注意的限制Zeta 引擎保证同一任务中不同连接器的 jar 分开加载Spark/Flink 引擎仍会将所有连接器依赖 jar 一起加载同一任务放置不同版本的 jar 可能导致冲突。此外任何不以connector-开头的目录或 jar 会被当作通用依赖处理Zeta 引擎中可将共享 jar 放入${SEATUNNEL_HOME}/lib/。六、选型总结按照本文的决策路径操作用第一节的四项清单锁定候选 Sink目标系统形态 → 写入保证 → 依赖 → 写入模式支持在 docs/zh/connectors/sink/ 找到对应连接器页面核对“主要特性”勾选框与参数表数据库类目标二选一定写入模式——generate_sink_sql true配合database/table/primary_keys可叠加 save mode或query完全自控 SQL跳过 save mode文件类目标先确认该连接器是否暴露schema_save_mode/data_save_mode再决定写入前的目录/数据处理策略多流作业务必显式配置plugin_input/plugin_output驱动与 SDK 按依赖隔离机制放入对应插件目录。更多边界情况与常见问题可继续参考 连接器常见问题 与 Sink 写入模式与 Save Mode。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表