
SeaTunnel Transforms 全景指南从数据接线到多表路由的字段级加工实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 Transform转换层位于 Source 与 Sink 之间是承担字段映射、行过滤、SQL 加工、表路由等管道中间逻辑的核心环节。本文以 SeaTunnel 官方 Transforms 总览为主线系统讲解 Transform 在整个作业中的位置、plugin_input/plugin_output数据集接线机制、常用转换插件的配置实战以及多表场景下的表路由与合并方案帮助你快速定位合适的转换插件并搭建出可读、可验证的数据管道。Transform 在作业中的位置与职责SeaTunnel 将一条数据管道抽象为三个逻辑阶段Source - Transform Chain - SinkTransform 块在作业配置中是可选的但当出现以下需求时它就是表达管道逻辑的主要场所Source 字段与 Sink 字段不能直接对齐需要字段映射、重命名或字段裁剪数据行需要被过滤、补全或重塑CDC 元数据需要被转换成下游友好的形态一个作业需要路由或重塑多张逻辑表。从源码结构看SeaTunnel 的 Transform 生态集中在 seatunnel-transforms-v2 模块中按功能划分出fieldmapper、filter、rename、split、sql、table、copy、jsonpath、metadata、rowkind、regexextract、replace、encrypt、chunk、nlpmodelLLM / Embedding、dynamiccompile、python等实现包覆盖了从基础字段操作到 AI 模型调用的全谱系能力。从系统层面看Transform 的职责不仅是字段级映射还包括在不绑定特定引擎记录类型的前提下重塑数据行、在增删改列时同步维护 Schema 信息、将 row kind 或事件时间等元数据暴露为普通字段、在多表作业中路由/合并/过滤逻辑表并保持作业逻辑的声明式表达以便同一管道跑在不同引擎上。更完整的体系化说明可参考 Transform Plugin System。快速选型从目标出发挑选 Transform官方总览文档为不同目标给出了明确起点以下表格保持原样并补全了目标对应的插件定位Goal目标Start here起点理解 Transform 如何连接数据集Transform Common Options过滤行或裁剪字段Filter 和 Field Mapper使用 SQL 风格表达式SQL 和 SQL Functions重命名或重塑字段Field Rename 和 Split处理多张表Transform Multi Table 和 Table MergeSeaTunnel 的 Transform 插件总体上可以归入几大类便于从功能性质上快速判断行投影与映射FieldMapper、FieldRename、Copy用于对齐源字段与下游 Schema 预期过滤与路由Filter、TableFilter、TableMerge决定哪些记录或哪些表继续流向后续环节SQL 与表达式处理SQL、JsonPath、RegexExtract适合用声明式表达式表达转换逻辑元数据与 CDC 适配Metadata、RowKindExtractor、FilterRowKind在 CDC 管道中尤为重要用于保留或重塑变更语义可编程或 AI 处理DynamicCompile、Python、LLM、Embedding用于需要外部模型、富计算或自定义业务逻辑的行处理。数据集接线plugin_input 与 plugin_output 详解SeaTunnel 的 Transform 共享一套极简的接线wiring选项。这些选项不定义转换逻辑本身而是定义 Transform 如何连接作业内的上游与下游数据集。废弃的旧选项名:::caution 注意旧的选项名source_table_name和result_table_name已废弃新配置请使用plugin_input和plugin_output。:::共享接线选项选项含义典型用途plugin_input声明当前 Transform 消费哪个上游数据集。省略时按配置顺序读取上一个插件的输出。需要从某个命名中间数据集读取、或作业不是简单线性链时使用。plugin_output将当前 Transform 结果注册为命名数据集供后续 Transform 或 Sink 引用。多个下游步骤需要同一结果、或想让管道图更显式时使用。两种接线模式数据集的接线在高层面上分为两种模式隐式链式implicit chaining每个插件按配置顺序读取上一个插件的输出写法最简短适合非常小的作业显式数据集接线explicit dataset wiring插件通过plugin_input与plugin_output引用命名数据集。显式命名在以下场景中更值得采用一个 Source 喂给多个下游步骤一个 Transform 的结果被多个 Sink 复用作业包含多张逻辑表希望管道图更易读、更易调试。源码佐证接线选项的底层定义接线选项在源码 TransformCommonOptions.java 中有直接体现例如table_transform、table_path、table_match_regex默认.*匹配所有表、rule_match_modeFIRST_MATCH/ALL_MATCH等选项都定义于此。此外还定义了row_error_handle_way默认fail可选skip、route_to_table与column_error_handle_way等错误处理选项说明 Transform 层不仅负责转换还内建了数据质量兜底策略——fail时格式错误会阻塞并抛异常skip时跳过该行数据ROUTE_TO_TABLE时可通过row_error_handle_way.error_table将脏数据路由到指定表。命名数据集流完整示例以下示例展示Source 注册数据集fake一个 Transform 读取该数据集并产出fake1两个 Sink 分别消费不同输出env { job.mode BATCH } source { FakeSource { plugin_output fake row.num 100 schema { fields { id int name string age int c_timestamp timestamp } } } } transform { Sql { plugin_input fake plugin_output fake1 query select id, upper(name) as name, age 1 as age, c_timestamp from fake } } sink { Console { plugin_input fake1 } Console { plugin_input fake } }实践准则数据集命名保持简短且有意义一旦作业出现分支或多表行为优先使用显式命名保持 Transform 文档与配置示例中的数据集名一致不要用数据集接线来掩盖过于复杂的逻辑当作业难以读懂时应考虑拆分作业。新手推荐学习路径官方总览文档给出的新用户建议路径分三步本文在此基础上补充了每个环节应掌握的关键点先读 Common Options 页把plugin_input和plugin_output的语义彻底弄清——这是理解后续所有示例的前提先选择与目标最匹配的最简 Transform再进入 SQL 或多表编排。例如只做字段裁剪就先用Filter只做改名就用FieldRename避免一上来就用 SQL 包揽一切逐步添加 Transform保持管道可读在验证作业的过程中每次只增加一步转换方便定位问题。常用转换插件实战速览以下四个插件是官方总览选型表中反复出现的核心成员均支持plugin_input/plugin_output公共选项详见 Common Options。Filter保留或剔除字段Filter提供include_fields与exclude_fields两个互斥选项必须且只能设置其中一个。include_fields列出需要保留的字段未列出的字段被删除exclude_fields列出需要删除的字段未列出的字段被保留。源数据nameagecardJoy Ding20123May Ding20123Kin Dom20123Joy Dom20123保留name、card字段transform { Filter { plugin_input fake plugin_output fake1 include_fields [name, card] } }或删除age字段transform { Filter { plugin_input fake plugin_output fake1 exclude_fields [age] } }当一个大表字段非常多、只需删除少量字段时exclude_fields特别实用。处理后结果表fake1如下namecardJoy Ding123May Ding123Kin Dom123Joy Dom123FieldMapper映射输入输出字段FieldMapper通过必填的field_mapper对象指定输入与输出之间的字段映射关系。示例中我们希望删除age字段、调整字段顺序为id、card、name并把name重命名为new_nametransform { FieldMapper { plugin_input fake plugin_output fake1 field_mapper { id id card card name new_name } } }映射后结果表fake1idcardnew_name1123Joy Ding2123May Ding3123Kin Dom4123Joy DomFieldRename批量重命名字段FieldRename支持多种重命名策略各选项均非必填可按需组合选项类型默认值说明convert_casestring-大小写转换类型取值为UPPER、LOWERprefixstring-追加到字段名前的前缀suffixstring-追加到字段名后的后缀replacements_with_regexarray-替换规则数组每条规则为replace_from、replace_to及可选is_regex默认trueis_regexfalse时按精确字段名全匹配处理specificarray-定向重命名规则每条为field_name与target_name命中后直接重命名并跳过其他规则在 CDC 场景下把 MySQL 大小写混合的字段名统一为大写并加前后缀env { parallelism 1 job.mode STREAMING } source { MySQL-CDC { plugin_output customers_mysql_cdc username root password 123456 table-names [source.user_shop, source.user_order] url jdbc:mysql://localhost:3306/source } } transform { FieldRename { plugin_input customers_mysql_cdc plugin_output trans_result convert_case UPPER prefix F_ suffix _S replacements_with_regex [ { replace_from create_time replace_to SOURCE_CREATE_TIME } ] } } sink { Jdbc { plugin_input trans_result driveroracle.jdbc.OracleDriver urljdbc:oracle:thin:oracle-host:1521/ORCLCDB usermyuser passwordmypwd generate_sink_sql true database ORCLCDB table ${database_name}.${table_name} primary_keys [${primary_key}] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }定向重命名单个字段如把InvoiceNum规范为invoice_numtransform { FieldRename { plugin_input input plugin_output output specific [ { field_name InvoiceNum, target_name invoice_num } ] } }Split按分隔符拆分为多字段Split将单个字段拆分为多个字段必填项为separator分隔符、split_field被拆字段、output_fields拆分后的结果字段数组。将name拆成first_name和last_nametransform { Split { plugin_input fake plugin_output fake1 separator split_field name output_fields [first_name, last_name] } }拆分后结果表fake1nameagecardfirst_namelast_nameJoy Ding20123JoyDingMay Ding20123MayDingKin Dom20123KinDomJoy Dom20123JoyDomSQL Transform声明式行级加工当转换逻辑适合用 SQL 表达时SqlTransform 是首选。它使用内存 SQL 引擎可借助 SQL 函数与引擎能力实现转换任务。核心选项选项类型必填默认值说明plugin_inputstring是-源表名query SQL 中的表名必须与之匹配plugin_outputstring是-输出数据集名querystring是-查询 SQL支持基础函数与条件过滤尚不支持多源表 JOIN 与聚合等复杂 SQLenginestring否ZETASQL 引擎支持ZETA与INTERNALquery既可以用select [table_name.]column_a查询普通列表名可省略也可以用select c_row.c_inner_row.column_b查询内嵌 struct 列此时表达式不允许出现表名。基础示例源数据idnameage1Joy Ding202May Ding213Kin Dom244Joy Dom22transform { Sql { plugin_input fake plugin_output fake1 query select id, concat(name, _) as name, age1 as age from dual where id0 } }结果表fake1idnameage1Joy Ding_212May Ding_223Kin Dom_254Joy Dom_23Struct 嵌套查询若上游 Schema 含嵌套结构source { FakeSource { plugin_output fake row.num 100 string.template [innerQuery] schema { fields { name string c_date date c_row { c_inner_row { c_inner_int int c_inner_string string c_inner_timestamp timestamp c_map_1 mapstring, string c_map_2 mapstring, mapstring,string } c_string string } } } } }以下查询均合法select name, c_date, c_row, c_row.c_inner_row, c_row.c_string, c_row.c_inner_row.c_inner_int, c_row.c_inner_row.c_inner_string, c_row.c_inner_row.c_inner_timestamp, c_row.c_inner_row.c_map_1, c_row.c_inner_row.c_map_1.some_key以下查询不合法——map 必须是最后一个结构不能查询嵌套的 mapselect c_row.c_inner_row.c_map_2.some_key.inner_map_keySQL 函数能力SqlTransform 内置了丰富的函数库覆盖字符串、数值、时间日期、系统函数、向量函数等类别完整清单见 SQL Functions。几个代表性的例子字符串CONCAT(NAME, _)、UPPER(NAME)、REGEXP_REPLACE(Hello World, , )、SPLIT(test, ;)、SUBSTRING([Hello], 2)数值ABS(I)、ROUND(N, 2)、POWER(A, B)、MOD(A, B)、CEIL(A)时间日期CURRENT_TIMESTAMP、DATEADD(CREATED, 1, MONTH)、DATE_TRUNC(CREATED, DAY)、EXTRACT(YEAR FROM TIMESTAMP 2001-02-16 20:38:40)、FROM_UNIXTIME(1672502400, yyyy-MM-dd HH:mm:ss, UTC6)、TO_DATE(2021-04-08, yyyy-MM-dd)系统函数CAST(NAME AS INT)、TRY_CAST(NAME AS INT)失败返回 NULL、COALESCE(A, B, C)、CASE WHEN ... THEN ... ELSE ... END、UUID()、ARRAY(1,2,3)、LATERAL VIEW EXPLODE(SPLIT(NAME, ,))向量函数VECTOR_DIMS(vector)、VECTOR_NORM(vector)、INNER_PRODUCT(v1, v2)、COSINE_DISTANCE(v1, v2)、L1_DISTANCE、L2_DISTANCE、VECTOR_REDUCE(embedding, 256, TRUNCATE)、VECTOR_NORMALIZE(embedding)。SQL Transform 的完整作业示例含 env、source、transform、sink 全链路可参考 SQL 中的 Job Config Example。多表 Transform一次配置处理多张表SeaTunnel 的 Transform 支持多表转换特别适合上游插件输出多张表的场景如JDBCSource、MySQL-CDC。所有 Transform 均可按多表方式配置且多表模式对转换能力没有任何限制——其目的是将多张表的转换配置合并为一个 Transform 便于管理。多表配置属性名称类型必填默认值说明table_match_regexString否.*匹配需要转换的表名指上游真实表名而非plugin_output的正则表达式默认匹配所有表table_transformList否-针对单张表的规则列表为某表配置了table_transform规则后外层规则不再作用于该表table_transform优先级更高table_transform.table_pathString否-指定表路径格式为databaseName[.schemaName].tableName精确匹配rule_match_modeString否-控制多条table_transform条目命中同一精确table_path时的求值方式取值FIRST_MATCH与ALL_MATCH匹配逻辑示例假设上游有 5 张结构相同字段id、name、age的表test.abc、test.abcd、test.xyz、test.xyzxyz、test.www。使用CopyTransform 实现差异化复制前两张表把name复制为name1test.xyz复制为name2test.xyzxyz复制为name3test.www保持不变transform { Copy { plugin_input fake // 可选数据集名 plugin_output fake1 // 可选数据集名 table_match_regex test.a.* // 匹配 test.abc 和 test.abcd src_field name dest_field name1 table_transform [{ table_path test.xyz src_field name dest_field name2 }, { table_path test.xyzxyz src_field name dest_field name3 }] } }各表的最终输出结构test.abc、test.abcd→id | name | age | name1test.xyz→id | name | age | name2test.xyzxyz→id | name | age | name3test.www→id | name | age不做转换。每张表的配置优先级为table_transformtable_match_regex若某表没有任何规则命中则不进行转换。table_transform.table_path采用精确表路径匹配rule_match_mode仅在多条条目使用同一精确table_path时生效未配置时配置解析阶段会拒绝重复的精确table_pathFIRST_MATCH按声明顺序应用第一条命中规则ALL_MATCH按声明顺序应用全部命中规则前一条规则的输出作为同表下一条规则的输入。例如ALL_MATCH模式下可先复制name为name2再复制name2为name3transform { Copy { rule_match_mode ALL_MATCH table_transform [{ table_path test.xyz src_field name dest_field name2 }, { table_path test.xyz src_field name2 dest_field name3 }] } }输出结构为id | name | age | name2 | name3。多表模式更多细节见 Transform Multi Table。TableMerge分表合并TableMerge用于合并分库分表数据选项包括database新库名可选、schema新 schema 名可选、table新表名必填。将source.user_1、source.user_2等分表合并为user_db.user_allenv { parallelism 1 job.mode STREAMING } source { MySQL-CDC { plugin_output customers_mysql_cdc username root password 123456 table-names [source.user_1, source.user_2, source.shop] url jdbc:mysql://localhost:3306/source } } transform { TableMerge { plugin_input customers_mysql_cdc plugin_output trans_result table_match_regex source.user_.* database user_db table user_all } } sink { Jdbc { plugin_input trans_result drivercom.mysql.cj.jdbc.Driver urljdbc:mysql://localhost:3306/sink usermyuser passwordmypwd generate_sink_sql true database ${database_name} table ${table_name} primary_keys [${primary_key}] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }小结SeaTunnel 的 Transform 层以plugin_input/plugin_output数据集接线为核心支持从简单的字段裁剪Filter、字段映射FieldMapper、字段重命名FieldRename、字段拆分Split到声明式的 SQL 加工Sql SQL Functions再到多表作业的表级路由多表 Transform与分表合并TableMerge。新手建议先掌握接线选项从最简插件入手再逐步叠加 SQL 与多表编排如需系统性理解 Transform 的契约设计与执行机制可继续阅读 Transform Plugin System 与 Job Configuration Guide以及完整的 Transforms 目录。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考