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 -> SinkTransform 块在作业配置中是可选的,但当出现以下需求时,它就是表达管道逻辑的主要场所:
- Source 字段与 Sink 字段不能直接对齐,需要字段映射、重命名或字段裁剪;
- 数据行需要被过滤、补全或重塑;
- CDC 元数据需要被转换成下游友好的形态;
- 一个作业需要路由或重塑多张逻辑表。
从源码结构看,SeaTunnel 的 Transform 生态集中在 seatunnel-transforms-v2 模块中,按功能划分出fieldmapper、filter、rename、split、sql、table、copy、jsonpath、metadata、rowkind、regexextract、replace、encrypt、chunk、nlpmodel(LLM / 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 Merge |
SeaTunnel 的 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_mode(FIRST_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列出需要删除的字段,未列出的字段被保留。
源数据:
| name | age | card |
|---|---|---|
| Joy Ding | 20 | 123 |
| May Ding | 20 | 123 |
| Kin Dom | 20 | 123 |
| Joy Dom | 20 | 123 |
保留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如下:
| name | card |
|---|---|
| Joy Ding | 123 |
| May Ding | 123 |
| Kin Dom | 123 |
| Joy Dom | 123 |
FieldMapper:映射输入输出字段
FieldMapper通过必填的field_mapper对象指定输入与输出之间的字段映射关系。示例中我们希望删除age字段、调整字段顺序为id、card、name,并把name重命名为new_name:
transform { FieldMapper { plugin_input = "fake" plugin_output = "fake1" field_mapper = { id = id card = card name = new_name } } }映射后结果表fake1:
| id | card | new_name |
|---|---|---|
| 1 | 123 | Joy Ding |
| 2 | 123 | May Ding |
| 3 | 123 | Kin Dom |
| 4 | 123 | Joy Dom |
FieldRename:批量重命名字段
FieldRename支持多种重命名策略,各选项均非必填,可按需组合:
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
convert_case | string | - | 大小写转换类型,取值为UPPER、LOWER |
prefix | string | - | 追加到字段名前的前缀 |
suffix | string | - | 追加到字段名后的后缀 |
replacements_with_regex | array | - | 替换规则数组,每条规则为replace_from、replace_to及可选is_regex(默认true);is_regex=false时按精确字段名(全匹配)处理 |
specific | array | - | 定向重命名规则,每条为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" 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_name和last_name:
transform { Split { plugin_input = "fake" plugin_output = "fake1" separator = " " split_field = "name" output_fields = [first_name, last_name] } }拆分后结果表fake1:
| name | age | card | first_name | last_name |
|---|---|---|---|---|
| Joy Ding | 20 | 123 | Joy | Ding |
| May Ding | 20 | 123 | May | Ding |
| Kin Dom | 20 | 123 | Kin | Dom |
| Joy Dom | 20 | 123 | Joy | Dom |
SQL Transform:声明式行级加工
当转换逻辑适合用 SQL 表达时,SqlTransform 是首选。它使用内存 SQL 引擎,可借助 SQL 函数与引擎能力实现转换任务。
核心选项
| 选项 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
plugin_input | string | 是 | - | 源表名,query SQL 中的表名必须与之匹配 |
plugin_output | string | 是 | - | 输出数据集名 |
query | string | 是 | - | 查询 SQL,支持基础函数与条件过滤;尚不支持多源表 JOIN 与聚合等复杂 SQL |
engine | string | 否 | ZETA | SQL 引擎,支持ZETA与INTERNAL |
query既可以用select [table_name.]column_a查询普通列(表名可省略),也可以用select c_row.c_inner_row.column_b查询内嵌 struct 列(此时表达式不允许出现表名)。
基础示例
源数据:
| id | name | age |
|---|---|---|
| 1 | Joy Ding | 20 |
| 2 | May Ding | 21 |
| 3 | Kin Dom | 24 |
| 4 | Joy Dom | 22 |
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:
| id | name | age |
|---|---|---|
| 1 | Joy Ding_ | 21 |
| 2 | May Ding_ | 22 |
| 3 | Kin Dom_ | 25 |
| 4 | Joy 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_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', '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 ... 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_regex | String | 否 | .* | 匹配需要转换的表名(指上游真实表名而非plugin_output)的正则表达式,默认匹配所有表 |
table_transform | List | 否 | - | 针对单张表的规则列表;为某表配置了table_transform规则后,外层规则不再作用于该表,table_transform优先级更高 |
table_transform.table_path | String | 否 | - | 指定表路径,格式为databaseName[.schemaName].tableName,精确匹配 |
rule_match_mode | String | 否 | - | 控制多条table_transform条目命中同一精确table_path时的求值方式,取值FIRST_MATCH与ALL_MATCH |
匹配逻辑示例
假设上游有 5 张结构相同(字段id、name、age)的表:test.abc、test.abcd、test.xyz、test.xyzxyz、test.www。使用CopyTransform 实现差异化复制:前两张表把name复制为name1,test.xyz复制为name2,test.xyzxyz复制为name3,test.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 | name1;test.xyz→id | name | age | name2;test.xyzxyz→id | name | age | name3;test.www→id | name | age(不做转换)。
每张表的配置优先级为:table_transform>table_match_regex;若某表没有任何规则命中,则不进行转换。table_transform.table_path采用精确表路径匹配,rule_match_mode仅在多条条目使用同一精确table_path时生效:未配置时配置解析阶段会拒绝重复的精确table_path;FIRST_MATCH按声明顺序应用第一条命中规则;ALL_MATCH按声明顺序应用全部命中规则,前一条规则的输出作为同表下一条规则的输入。例如ALL_MATCH模式下可先复制name为name2,再复制name2为name3:
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_1、source.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),仅供参考