SeaTunnel Transforms 全景指南:从数据接线到多表路由的字段级加工实战
2026/9/18 7:31:22 网站建设 项目流程

SeaTunnel Transforms 全景指南:从数据接线到多表路由的字段级加工实战

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

SeaTunnel 的 Transform(转换)层位于 Source 与 Sink 之间,是承担字段映射、行过滤、SQL 加工、表路由等管道中间逻辑的核心环节。本文以 SeaTunnel 官方 Transforms 总览为主线,系统讲解 Transform 在整个作业中的位置、plugin_input/plugin_output数据集接线机制、常用转换插件的配置实战,以及多表场景下的表路由与合并方案,帮助你快速定位合适的转换插件并搭建出可读、可验证的数据管道。

Transform 在作业中的位置与职责

SeaTunnel 将一条数据管道抽象为三个逻辑阶段:

Source -> Transform Chain -> Sink

Transform 块在作业配置中是可选的,但当出现以下需求时,它就是表达管道逻辑的主要场所:

  • Source 字段与 Sink 字段不能直接对齐,需要字段映射、重命名或字段裁剪;
  • 数据行需要被过滤、补全或重塑;
  • CDC 元数据需要被转换成下游友好的形态;
  • 一个作业需要路由或重塑多张逻辑表。

从源码结构看,SeaTunnel 的 Transform 生态集中在 seatunnel-transforms-v2 模块中,按功能划分出fieldmapperfilterrenamesplitsqltablecopyjsonpathmetadatarowkindregexextractreplaceencryptchunknlpmodel(LLM / Embedding)、dynamiccompilepython等实现包,覆盖了从基础字段操作到 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 Merge

SeaTunnel 的 Transform 插件总体上可以归入几大类,便于从功能性质上快速判断:

  • 行投影与映射FieldMapperFieldRenameCopy,用于对齐源字段与下游 Schema 预期;
  • 过滤与路由FilterTableFilterTableMerge,决定哪些记录或哪些表继续流向后续环节;
  • SQL 与表达式处理SQLJsonPathRegexExtract,适合用声明式表达式表达转换逻辑;
  • 元数据与 CDC 适配MetadataRowKindExtractorFilterRowKind,在 CDC 管道中尤为重要,用于保留或重塑变更语义;
  • 可编程或 AI 处理DynamicCompilePythonLLMEmbedding,用于需要外部模型、富计算或自定义业务逻辑的行处理。

数据集接线:plugin_input 与 plugin_output 详解

SeaTunnel 的 Transform 共享一套极简的接线(wiring)选项。这些选项不定义转换逻辑本身,而是定义 Transform 如何连接作业内的上游与下游数据集。

废弃的旧选项名

:::caution 注意

旧的选项名source_table_nameresult_table_name已废弃,新配置请使用plugin_inputplugin_output

:::

共享接线选项

选项含义典型用途
plugin_input声明当前 Transform 消费哪个上游数据集。省略时,按配置顺序读取上一个插件的输出。需要从某个命名中间数据集读取、或作业不是简单线性链时使用。
plugin_output将当前 Transform 结果注册为命名数据集,供后续 Transform 或 Sink 引用。多个下游步骤需要同一结果、或想让管道图更显式时使用。

两种接线模式

数据集的接线在高层面上分为两种模式:

  • 隐式链式(implicit chaining):每个插件按配置顺序读取上一个插件的输出,写法最简短,适合非常小的作业;
  • 显式数据集接线(explicit dataset wiring):插件通过plugin_inputplugin_output引用命名数据集。

显式命名在以下场景中更值得采用:

  • 一个 Source 喂给多个下游步骤;
  • 一个 Transform 的结果被多个 Sink 复用;
  • 作业包含多张逻辑表;
  • 希望管道图更易读、更易调试。

源码佐证:接线选项的底层定义

接线选项在源码 TransformCommonOptions.java 中有直接体现,例如table_transformtable_pathtable_match_regex(默认.*,匹配所有表)、rule_match_modeFIRST_MATCH/ALL_MATCH)等选项都定义于此。此外还定义了row_error_handle_way(默认fail,可选skiproute_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 文档与配置示例中的数据集名一致;
  • 不要用数据集接线来掩盖过于复杂的逻辑,当作业难以读懂时应考虑拆分作业。

新手推荐学习路径

官方总览文档给出的新用户建议路径分三步,本文在此基础上补充了每个环节应掌握的关键点:

  1. 先读 Common Options 页,把plugin_inputplugin_output的语义彻底弄清——这是理解后续所有示例的前提;
  2. 先选择与目标最匹配的最简 Transform,再进入 SQL 或多表编排。例如只做字段裁剪就先用Filter,只做改名就用FieldRename,避免一上来就用 SQL 包揽一切;
  3. 逐步添加 Transform,保持管道可读,在验证作业的过程中每次只增加一步转换,方便定位问题。

常用转换插件实战速览

以下四个插件是官方总览选型表中反复出现的核心成员,均支持plugin_input/plugin_output公共选项(详见 Common Options)。

Filter:保留或剔除字段

Filter提供include_fieldsexclude_fields两个互斥选项(必须且只能设置其中一个)。include_fields列出需要保留的字段,未列出的字段被删除;exclude_fields列出需要删除的字段,未列出的字段被保留。

源数据:

nameagecard
Joy Ding20123
May Ding20123
Kin Dom20123
Joy Dom20123

保留namecard字段:

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如下:

namecard
Joy Ding123
May Ding123
Kin Dom123
Joy Dom123

FieldMapper:映射输入输出字段

FieldMapper通过必填的field_mapper对象指定输入与输出之间的字段映射关系。示例中我们希望删除age字段、调整字段顺序为idcardname,并把name重命名为new_name

transform { FieldMapper { plugin_input = "fake" plugin_output = "fake1" field_mapper = { id = id card = card name = new_name } } }

映射后结果表fake1

idcardnew_name
1123Joy Ding
2123May Ding
3123Kin Dom
4123Joy Dom

FieldRename:批量重命名字段

FieldRename支持多种重命名策略,各选项均非必填,可按需组合:

选项类型默认值说明
convert_casestring-大小写转换类型,取值为UPPERLOWER
prefixstring-追加到字段名前的前缀
suffixstring-追加到字段名后的后缀
replacements_with_regexarray-替换规则数组,每条规则为replace_fromreplace_to及可选is_regex(默认true);is_regex=false时按精确字段名(全匹配)处理
specificarray-定向重命名规则,每条为field_nametarget_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" driver="oracle.jdbc.OracleDriver" url="jdbc:oracle:thin:@oracle-host:1521/ORCLCDB" user="myuser" password="mypwd" 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_num):

transform { FieldRename { plugin_input = "input" plugin_output = "output" specific = [ { field_name = "InvoiceNum", target_name = "invoice_num" } ] } }

Split:按分隔符拆分为多字段

Split将单个字段拆分为多个字段,必填项为separator(分隔符)、split_field(被拆字段)、output_fields(拆分后的结果字段数组)。将name拆成first_namelast_name

transform { Split { plugin_input = "fake" plugin_output = "fake1" separator = " " split_field = "name" output_fields = [first_name, last_name] } }

拆分后结果表fake1

nameagecardfirst_namelast_name
Joy Ding20123JoyDing
May Ding20123MayDing
Kin Dom20123KinDom
Joy Dom20123JoyDom

SQL Transform:声明式行级加工

当转换逻辑适合用 SQL 表达时,SqlTransform 是首选。它使用内存 SQL 引擎,可借助 SQL 函数与引擎能力实现转换任务。

核心选项

选项类型必填默认值说明
plugin_inputstring-源表名,query SQL 中的表名必须与之匹配
plugin_outputstring-输出数据集名
querystring-查询 SQL,支持基础函数与条件过滤;尚不支持多源表 JOIN 与聚合等复杂 SQL
enginestringZETASQL 引擎,支持ZETAINTERNAL

query既可以用select [table_name.]column_a查询普通列(表名可省略),也可以用select c_row.c_inner_row.column_b查询内嵌 struct 列(此时表达式不允许出现表名)。

基础示例

源数据:

idnameage
1Joy Ding20
2May Ding21
3Kin Dom24
4Joy Dom22
transform { Sql { plugin_input = "fake" plugin_output = "fake1" query = "select id, concat(name, '_') as name, age+1 as age from dual where id>0" } }

结果表fake1

idnameage
1Joy Ding_21
2May Ding_22
3Kin Dom_25
4Joy Dom_23

Struct 嵌套查询

若上游 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 = "map<string, string>" c_map_2 = "map<string, map<string,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 必须是最后一个结构,不能查询嵌套的 map:

select c_row.c_inner_row.c_map_2.some_key.inner_map_key

SQL 函数能力

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_TIMESTAMPDATEADD(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', 'UTC+6')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 ... ENDUUID()ARRAY(1,2,3)LATERAL VIEW EXPLODE(SPLIT(NAME, ','))
  • 向量函数VECTOR_DIMS(vector)VECTOR_NORM(vector)INNER_PRODUCT(v1, v2)COSINE_DISTANCE(v1, v2)L1_DISTANCEL2_DISTANCEVECTOR_REDUCE(embedding, 256, 'TRUNCATE')VECTOR_NORMALIZE(embedding)

SQL Transform 的完整作业示例(含 env、source、transform、sink 全链路)可参考 SQL 中的 Job Config Example。

多表 Transform:一次配置处理多张表

SeaTunnel 的 Transform 支持多表转换,特别适合上游插件输出多张表的场景(如JDBCSourceMySQL-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_MATCHALL_MATCH

匹配逻辑示例

假设上游有 5 张结构相同(字段idnameage)的表:test.abctest.abcdtest.xyztest.xyzxyztest.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.abctest.abcdid | name | age | name1
  • test.xyzid | name | age | name2
  • test.xyzxyzid | name | age | name3
  • test.wwwid | name | age(不做转换)。

每张表的配置优先级为:table_transform>table_match_regex;若某表没有任何规则命中,则不进行转换。table_transform.table_path采用精确表路径匹配,rule_match_mode仅在多条条目使用同一精确table_path时生效:未配置时配置解析阶段会拒绝重复的精确table_pathFIRST_MATCH按声明顺序应用第一条命中规则;ALL_MATCH按声明顺序应用全部命中规则,前一条规则的输出作为同表下一条规则的输入。例如ALL_MATCH模式下可先复制namename2,再复制name2name3

transform { 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_1source.user_2等分表合并为user_db.user_all

env { 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" driver="com.mysql.cj.jdbc.Driver" url="jdbc:mysql://localhost:3306/sink" user="myuser" password="mypwd" 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询