ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

SeaTunnel Sink 写入模式与 Save Mode:generate_sink_sql、query、schema_save_mode 与 data_save_mode 配置实战

SeaTunnel Sink 写入模式与 Save Mode:generate_sink_sql、query、schema_save_mode 与 data_save_mode 配置实战 SeaTunnel Sink 写入模式与 Save Modegenerate_sink_sql、query、schema_save_mode 与 data_save_mode 配置实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文围绕 SeaTunnel 中 Sink 配置最容易混淆的两个决策展开写入模式SeaTunnel 如何把每一行数据写到目标端与 Save Mode开始写入前如何处理目标端已存在的表、索引、目录或数据。读完本文你将能够在generate_sink_sql、query、schema_save_mode、data_save_mode、custom_sql、primary_keys、enable_upsert之间做出正确选择理解各参数在 JdbcSinkFactory 与 DefaultSaveModeHandler 中的实际执行逻辑并掌握 JDBC 与文件类 Sink 的参数边界和常见故障排查方法。两个核心概念写入模式与 Save Mode在 SeaTunnel 的 Sink 配置里有两组经常被混用的参数写入模式决定 SeaTunnel 如何把每一行数据写到目标端。对 JDBC 目标端而言就是“让 SeaTunnel 自动生成 INSERT / UPSERT / UPDATE / DELETE SQL”还是“我完全控制写入 SQL”Save Mode决定 SeaTunnel 在开始写入数据前如何处理目标端已经存在的表、索引、目录或数据拆分为schema_save_mode结构层面和data_save_mode数据层面两个参数。从源码结构看Save Mode 的语义在 API 层统一实现DefaultSaveModeHandler 中handleSchemaSaveMode()与handleDataSaveMode()两个方法分别对SchemaSaveMode和DataSaveMode枚举做 switch 分发各 connector如 JdbcSaveModeHandler继承它并按需覆写建表等行为。快速决策表目标优先选择说明让 SeaTunnel 为 JDBC 目标端生成 INSERT / UPSERT / UPDATE / DELETE SQLgenerate_sink_sql true并配置database、table通常还要配置primary_keys这是 JDBC Sink 的能力。SeaTunnel 能解析目标 Catalog 表时也可以执行 save mode 和自动建表。完全控制 JDBC 写入 SQLquery INSERT ... VALUES (?, ...)不要和generate_sink_sql true同时配置。JDBC Sink 在这个模式下不会执行schema_save_mode、data_save_mode或custom_sql。目标表不存在时自动创建或不存在时报错schema_save_mode仅适用于显式暴露 save mode 参数并且能通过 Catalog 创建或检查目标端的 Sink。写入前保留、清空或检查目标端已有数据data_save_mode支持值取决于具体 connector。File Sink 通常只支持DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS。写入数据前先执行一条自定义 SQLdata_save_mode CUSTOM_PROCESSING和custom_sql仅适用于同时暴露这两个参数的 connector。这是写入前钩子不是逐行写入 SQL。JDBC Sink 使用数据库原生 Upsertgenerate_sink_sql true、primary_keys、enable_upsert true没有可用主键或唯一键时JDBC 自动生成 SQL 会退化为普通 INSERT。写入对象存储或文件系统查看具体 File Sink 参数表File Sink 不使用generate_sink_sql。部分文件 connector 暴露 save mode部分不暴露。JDBCgenerate_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和query。这一点不只是文档约定而是由 OptionRule 在提交期强制校验的。在 JdbcSinkFactory.optionRule() 中可以看到条件约束.conditional(JdbcSinkOptions.GENERATE_SINK_SQL, true, JdbcSinkOptions.DATABASE) .conditional(JdbcSinkOptions.GENERATE_SINK_SQL, false, JdbcSinkOptions.QUERY) .conditional( JdbcSinkOptions.DATA_SAVE_MODE, DataSaveMode.CUSTOM_PROCESSING, JdbcSinkOptions.CUSTOM_SQL)即generate_sink_sql true时database变为必填generate_sink_sql false时要求配置querydata_save_mode CUSTOM_PROCESSING时custom_sql变为必填。再看默认值见 JdbcSinkOptions参数默认值说明generate_sink_sqlfalse即默认走自定义query模式需要自动生成 SQL 时必须显式开启schema_save_modeCREATE_SCHEMA_WHEN_NOT_EXIST表不存在时自动建表data_save_modeAPPEND_DATA保留已有数据并追加写入enable_upserttrue拿到可用主键/唯一键时启用原生 upsertcreate_indextrue自动建表时是否创建索引此外schema_save_mode与data_save_mode在 OptionRule 中属于required项有默认值但仍参与校验链路custom_sql只在CUSTOM_PROCESSING场景下强制要求。主键是如何被解析的使用generate_sink_sql true时如果目标端需要处理 UPDATE、DELETE 或 upsert 记录请配置primary_keys。如果没有显式配置SeaTunnel 会尝试从上游 Catalog 元数据继承主键再尝试第一组唯一键仍然没有可用键时会退化为普通 INSERT。这一继承逻辑在 JdbcSinkFactory.createSink() 中可以清晰看到先读取上游catalogTable的PrimaryKey若配置里没有primary_keys且上游 schema 带了主键则把主键列名注入配置否则遍历getConstraintKeys()找第一组UNIQUE_KEY约束并注入。若显式配置了primary_keys则会用它构造新的PrimaryKey覆盖CatalogTable的 schemaL145-L163后续生成的 upsert/update/delete SQL 均以该主键为准。因此当 CDC 链路或需要精确 upsert 的场景下建议显式配置primary_keys避免依赖上游元数据带来的不确定性。Save Mode 语义详解schema_save_modeschema_save_mode控制写入前如何处理目标结构。值行为RECREATE_SCHEMA目标不存在时创建目标已存在时删除后重建。CREATE_SCHEMA_WHEN_NOT_EXIST仅在目标不存在时创建。ERROR_WHEN_SCHEMA_NOT_EXIST目标不存在时报错。IGNORE跳过结构处理。对于数据库 Sink目标通常是表对于文件 Sink目标通常是路径或目录。DefaultSaveModeHandler 中每种模式都有对应实现RECREATE_SCHEMA走recreateSchema()先dropTable再createTableCREATE_SCHEMA_WHEN_NOT_EXIST只在!tableExists()时建表ERROR_WHEN_SCHEMA_NOT_EXIST在表不存在时抛出SINK_TABLE_NOT_EXIST异常。对于 JDBC 而言JdbcSaveModeHandler 覆写了createTable()先做createTablePreCheck()父类中会检查并自动创建不存在的 database再调用catalog.createTable(tablePath, catalogTable, true, createIndex)其中createIndex即上文提到的create_index参数。data_save_modedata_save_mode控制写入前如何处理目标端已有数据。值行为DROP_DATA保留结构并清空已有数据对应实现keepSchemaDropData()对已存在且非新建的表执行 truncate。APPEND_DATA保留已有数据并追加写入keepSchemaAndData()为空实现即什么都不做。CUSTOM_PROCESSING写入前执行custom_sql。仅适用于同时暴露这两个参数的 connector。ERROR_WHEN_DATA_EXISTS发现已有数据时报错对应实现中若dataExists()为真则抛出异常。需要强调custom_sql是写入前钩子不是逐行写入 SQL并且它只有在data_save_mode CUSTOM_PROCESSING且 save mode 处理真正执行时才会执行。从源码看DefaultSaveModeHandler.customProcessing()直接调用executeCustomSql()而customProcessing()只会在handleDataSaveMode()分发到CUSTOM_PROCESSING分支时被触发。这些参数不是所有 connector 都支持。最终请以你正在使用版本的具体 connector 参数表为准。Connector 支持边界JDBC 系列 SinkJDBC Sink 以及 MySQL、PostgreSQL、Oracle、SQL Server 等 JDBC 系列 Sink 页面使用同一套 JDBC 写入模式支持generate_sink_sql支持query自动生成 SQL 模式下支持schema_save_mode和data_save_modecustom_sql只有在 save mode 处理真正执行时才会执行即必须同时满足“自动生成 SQL 模式”和data_save_mode CUSTOM_PROCESSINGenable_upsert只有在 SeaTunnel 拿到可用主键或唯一键后才有意义。完整参数和示例请看 JDBC Sink。Doris SinkDoris Sink 支持schema_save_mode、data_save_mode、custom_sql和save_mode_create_template但不使用 JDBC 的generate_sink_sql。如果要处理 CDC DELETE 事件还需要 Doris 侧支持删除能力并按场景配置 connector 的sink.enable-delete。详见 Doris Sink。File 与对象存储 SinkFile Sink 写的是文件因此不使用generate_sink_sql、query或数据库 upsert。当前不同文件 connector 的支持边界如下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_mode。CosFile否当前 sink option rule 没有暴露schema_save_mode或data_save_mode。BosFile否当前 sink option rule 没有暴露schema_save_mode或data_save_mode。如果某个文件 connector 页面没有列出schema_save_mode或data_save_mode不要默认认为该 connector 可以接收这些参数。配置示例示例一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 } }要点generate_sink_sql true触发 OptionRule 对database的必填校验table public.orders这种带点号的写法会被 resolveSinkTablePath() 解析为 schema table按.切分显式给出primary_keys [id]保证 CDC 的 UPDATE / DELETE 事件生成 upsert/update/delete SQL 而不是普通 INSERTschema_save_mode保证目标表不存在时自动创建自动建库逻辑也包含在createTablePreCheck()中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。JDBC 代码中也明确注释了这一点use query to write data can not support get catalog见 JdbcSink.java即自定义 query 模式下无法获取 Catalogsave mode 因此无从执行。示例三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 } }DROP_DATA会保留目录结构但清空其中的已有文件数据后再开始写入如果希望失败时保护已有数据可改用ERROR_WHEN_DATA_EXISTS在检测到已有数据时直接报错。故障排查配了generate_sink_sql true但仍然只是 INSERT检查 SeaTunnel 是否拿到了可用 key。回顾 JdbcSinkFactory.createSink() 的继承顺序显式primary_keys 上游 Catalog 主键 第一组唯一键 退化为普通 INSERT。需要 upsert、update 或 delete 时建议显式配置primary_keys这是最稳妥的做法。JDBC Sink 的custom_sql没有执行按以下顺序排查Sink 是否配置了query。JDBC 自定义 query 模式不会执行 save mode 处理因此会跳过custom_sqldata_save_mode是否为CUSTOM_PROCESSING。只有这个值才会分发到executeCustomSql()其他模式下custom_sql会被忽略OptionRule 要求CUSTOM_PROCESSING必须与custom_sql成对出现反过来只配custom_sql而不改data_save_mode也是无效的。File Sink 不接受data_save_mode检查具体 connector 参数表。S3File、OssFile、HdfsFile、FtpFile、SftpFile、LocalFile暴露文件 save mode 参数ObsFile、CosFile和BosFile当前未暴露。参数是否被接受取决于各 connector 自己的 OptionRule 声明配置了未暴露的参数不会自动生效。我只想建表不想抽取数据Save mode 是 Sink 作业写入前的一部分SeaTunnel 目前没有通过schema_save_mode提供独立的“只执行 DDL”模式。如果作业没有数据Sink 仍可能完成初始化但这不能替代专门的 schema 管理流程。小结场景推荐配置组合常规数据库写入 自动建表generate_sink_sql trueschema_save_mode CREATE_SCHEMA_WHEN_NOT_EXISTdata_save_mode APPEND_DATACDC 全量 增量含 update/delete上行同上且显式配置primary_keys、enable_upsert true写入前需执行自定义 SQL如清理分区data_save_mode CUSTOM_PROCESSINGcustom_sql且必须处于自动生成 SQL 模式必须完全控制写入 SQLquery INSERT ... VALUES (?, ...)并接受 Save Mode 全部跳过对象存储 / 文件系统全量覆盖对暴露 save mode 的文件 connector 使用data_save_mode DROP_DATA记住两条核心边界写入模式generate_sink_sql/query决定“每行怎么写”且二者互斥Save Modeschema_save_mode/data_save_mode决定“写之前怎么处理目标端”且只在自动生成 SQL 模式、且目标端可被 Catalog 解析时才会真正执行。具体 connector 的参数边界永远以该 connector 的参数表为准。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表