SeaTunnel SQL Transform 完整实战指南:用内存 SQL 引擎完成行级数据转换

发布时间:2026/9/18 9:36:36

SeaTunnel SQL Transform 完整实战指南:用内存 SQL 引擎完成行级数据转换 SeaTunnel SQL Transform 完整实战指南用内存 SQL 引擎完成行级数据转换【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文以 docs/en/transforms/sql.md 为核心骨架结合seatunnel-transforms-v2模块中org.apache.seatunnel.transform.sql包的源码实现系统讲解 SeaTunnel 内置Sql变换插件的配置方式、SQL 语法能力边界、嵌套结构Struct/Map查询规则以及底层执行原理。读完本文你将能够在 SeaTunnel 作业中直接以 SQL 表达式对上游数据行进行字段投影、函数计算与条件过滤并理解其内存 SQL 引擎的工作机制与适用限制。一、SQL Transform 是什么SQL Transform 是 SeaTunnel 提供的行级数据变换插件使用内存 SQL 引擎对流入的每一行数据执行 SQL 表达式从而完成转换任务。在 SeaTunnel 的 Transform 体系中它位于 Source 与 Sink 之间用于字段映射、过滤、SQL 处理等管道中间操作参见 Transforms Overview。从源码看该插件的入口类为 SQLTransform.java其PLUGIN_NAME Sql在配置文件中以Sql { ... }的形式声明。它继承自AbstractCatalogSupportFlatMapTransform意味着这是一个一进多出flatMap风格的变换一条输入行经过 SQL 处理后可能输出零条被WHERE过滤掉、一条或多条记录。SQL Transform 的特点是逐行独立处理它不维护跨行的状态因此天然无法支持多表 JOIN 与跨行聚合AGGREGATE等复杂 SQL 操作这也是本文后续会重点强调的语法边界。二、核心选项详解Sql变换插件的配置选项如下表所示名称类型是否必填默认值plugin_inputstring是-plugin_outputstring是-querystring是-enginestring否ZETAplugin_input [string]输入表名上游数据的表名。query中的 SQL 表名必须与plugin_input指定的名称匹配在 SQL 中常写作from dual的形式详见下文。plugin_input/plugin_output是所有 Transform 插件的公共选项Common Options用于声明该变换读取哪张表、输出到哪张表。从 SQLTransform.java 的源码可以看到如果配置了plugin_input则取列表中第一个标识符作为inputTableName如果未配置则回退使用catalogTable.getTableId().getTableName()即上游 Catalog 表名。query [string]要执行的查询 SQL。它是一条简单的 SQL支持基础函数与条件过滤操作但复杂 SQL 暂不支持包括多源表/多行的 JOIN、AGGREGATE聚合操作等。查询表达式的语法要点select [table_name.]column_a查询名为column_a的列表名前缀是可选的即既可以写select fake.id也可以写select idselect c_row.c_inner_row.column_b查询内嵌结构体中的字段——即c_row列内部的c_inner_row列内部的column_b字段。注意在这种嵌套查询表达式中不能带表名。engine [string]该变换使用的 SQL 引擎支持的取值有ZETA与INTERNAL。未配置时默认使用ZETA。这里有一个值得注意的源码事实查看 SQLEngineFactory.java 可以发现ZETA与INTERNAL两种枚举值在工厂方法中均返回同一个ZetaSQLEngine实例switch (engineType) { case ZETA: case INTERNAL: return new ZetaSQLEngine(); }也就是说在当前仓库版本中无论你配置ZETA还是INTERNAL实际执行的都是基于 Zeta SQL 引擎的实现而engine选项的作用在于为未来接入其他引擎如基于 Apache Calcite 的引擎预留扩展位。此外SQLTransform在读取配置时会将引擎名toUpperCase()后通过EngineType.valueOf解析因此大小写不敏感参见 SQLTransform.java。提示SeaTunnel 还有一个独立的Calcite Transform插件PLUGIN_NAME Calcite由 CalciteTransform.java 实现采用 Apache Calcite 编译并执行 SQL。它与本文讲解的Sql插件是两个不同的插件engine选项并不用于切换到这个 Calcite 插件。两者详细对比如下。三、快速上手完整作业配置示例下面是一个完整的、可直接运行的 BATCH 作业配置原文示例来源 sql.md 的 Job Config Exampleenv { job.mode BATCH } source { FakeSource { plugin_output fake row.num 100 schema { fields { id int name string age int } } } } transform { Sql { plugin_input fake plugin_output fake1 query select id, concat(name, _) as name, age1 as age from dual where id0 } } sink { Console { plugin_input fake1 } }该作业的流转过程FakeSource生成 100 行数据字段为idint、namestring、ageint输出到表fakeSql变换读取表fake执行select id, concat(name, _) as name, age1 as age from dual where id0将结果输出到表fake1ConsoleSink 读取表fake1并打印到控制台。转换效果演示假设上游数据表内容如下idnameage1Joy Ding202May Ding213Kin Dom244Joy Dom22执行上述查询后结果表fake1中的数据将更新为idnameage1Joy Ding_212May Ding_223Kin Dom_254Joy Dom_23可以看到name列经过concat(name, _)追加了下划线后缀age列经过age1整体加一where id0则保证了所有行都被保留若把条件改成id2则前两行会被过滤掉。四、query 语法能力与边界源码级验证query支持简单 SQL即基础函数、字段投影与WHERE过滤。为了精确理解边界我们来看 ZetaSQLEngine.java 中validateSQL方法的实现它在作业启动解析 SQL 时即对语法做了严格校验if (!(statement instanceof Select)) { throw new IllegalArgumentException(Only supported DQL(select) SQL); } // 不支持 schema 前缀、表别名 // 不支持子查询sub table syntax // 不支持 JOIN // 不支持 ORDER BY // 不支持 GROUP BY // 不支持 LIMIT / OFFSET由此可以总结出query的完整能力与禁用清单支持SELECT查询DQL包括select *全列投影与显式列投影单表FROM子句表名必须与plugin_input或上游 Catalog 表名一致也可以使用dual作为占位表名基础内置函数调用字符串、数值、日期时间、系统函数等详见下文算术表达式如age1与别名asWHERE条件过滤、、、and/or等反引号包裹的字段名引擎会通过cleanEscape去除转义符参见 ZetaSQLEngine.java。不支持作业启动时会直接抛错JOIN多表关联GROUP BY与聚合操作SUM、COUNT、AVG等ORDER BY排序LIMIT/OFFSET分页子查询sub query带 schema 前缀的表名与表别名table alias非SELECT语句如INSERT、UPDATE、DELETE、DDL。提示如果表名与输入表名不一致引擎不会直接抛错而是打印一条 warn 日志SQL table name ... is not equal to input table name ...除非是DUAL这一点在排查问题时值得留意。支持的内置函数query中可以调用的内置函数非常丰富完整清单见 SQL Functions 文档按类别包括字符串函数CONCAT、CONCAT_WS、LOWER/UPPER、SUBSTRING/SUBSTR、TRIM/LTRIM/RTRIM、LPAD/RPAD、REPLACE、REGEXP_REPLACE、REGEXP_LIKE、REGEXP_SUBSTR、SPLIT、TO_CHAR等数值函数ABS、CEIL/CEILING、FLOOR、ROUND、MOD、EXP、LN、LOG/LOG10、SQRT、POWER、SIN/COS/TAN系列等时间与日期函数FORMATDATETIME、CURRENT_TIMESTAMP、日期加减与提取等系统函数CAST类型转换、COALESCE、IF等向量函数在支持向量类型的场景下可用。其中FORMATDATETIME(create_time,yyyy-MM-dd HH:mm)这类日期格式化函数在生成分区键等场景非常常用SQLTransformTest.java 的测试用例中就给出了该写法的示例。类型与 Schema 推导query的输出字段类型由引擎自动推导ZetaSQLEngine.typeMapping会根据SELECT项逐个推导输出列名与数据类型。若SELECT项是普通列且无别名输出列名沿用原列名若有别名则使用别名若为表达式则使用表达式字符串作为列名。推导结果会同步保留上游列的精度信息——测试用例testScaleSupport验证了时间戳列、字符串列的scale/columnLength在变换后正确保留见 SQLTransformTest.java。如果上游主键列全部出现在输出列中变换后的表会继承主键定义约束键ConstraintKey同理只有其涉及的全部列仍存在于输出中才会被继承参见 SQLTransform.java。五、嵌套结构查询Struct QuerySQL Transform 支持对 SeaTunnel 的复合类型嵌套结构体 Struct、Map进行字段级查询。上游 Schema 示例假设上游FakeSource的数据 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 } } } } }这是一个典型的三层嵌套结构顶层字段c_row是结构体内部包含c_inner_row再次嵌套结构体与c_string最内层又包含基本类型字段与两个 Map 字段c_map_1为普通 Mapc_map_2为 Map 套 Map 的二级嵌套 Map。合法的嵌套查询以下查询全部合法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要点拆解select name、select c_date直接投影顶层基本类型字段select c_row整体投影整个结构体列select c_row.c_inner_row投影结构体中的子结构体select c_row.c_inner_row.c_inner_int以点号逐层下钻访问最内层基本字段select c_row.c_inner_row.c_map_1.some_key可以读取 Map 中指定 key 的值c_map_1的 key 为some_key。不合法的查询以下查询不合法select c_row.c_inner_row.c_map_2.some_key.inner_map_key原因c_map_2是mapstring, mapstring,string类型的二级嵌套 Map而引擎要求Map 必须是查询路径上最后出现的结构the map must be the latest struct即不能在 Map 之后继续下钻。也就是说你可以查询c_map_2整体、可以查询c_map_2的某个 key得到mapstring,string类型但不能像上面那样对c_map_2的 key 再次取 key。这一限制同样在 ZetaSQLType.java 的类型推导逻辑中体现。嵌套查询的使用注意在前述嵌套查询表达式不能带表名的规则下c_row.c_inner_row.c_inner_int这类写法前面不能加表名前缀不能写成fake.c_row.c_inner_row.c_inner_intStruct 查询对 CDC、JSON 嵌套等复杂数据源的字段下钻十分有用配合COPY、FIELD_MAPPER等变换可以实现细粒度的结构重组。六、底层执行原理Scan → Filter → Project了解Sql变换的底层执行链路有助于写出更高效的 query。从 ZetaSQLEngine.java 的transformBySQL实现看Zeta 引擎对每一行数据执行的是一个经典的物理查询计划Scan Table扫描将输入行SeaTunnelRow的字段数组直接取出inputRow.getFields()作为后续计算的输入Filter过滤调用zetaSQLFilter.executeFilter(selectBody.getWhere(), inputFields)执行WHERE条件。若条件不满足返回 false则该方法返回null表示当前行被过滤、不产生任何输出行若WHERE表达式执行出错会抛出sqlWhereStatementErrorProject投影对保留下来的行遍历SELECT列表逐项计算。遇到AllColumnsselect *时原样展开所有输入字段遇到普通表达式时调用zetaSQLFunction.computeForValue(expression, inputFields)计算表达式的值函数调用、算术运算、嵌套字段访问等都在这一步完成表达式计算异常会抛出sqlExpressionError构造输出行将投影结果封装为新的SeaTunnelRow并保留输入行的 RowKind、TableId、Options 等元数据这对 CDC 场景下区分 INSERT/UPDATE/DELETE 语义至关重要Lateral View可选如果 SQL 中使用了LATERAL VIEW展开集合类型引擎会在此阶段把单行扩展为多行输出这也是transformBySQL返回ListSeaTunnelRow而不是单行的原因之一。引擎生命周期SQLTransform.open()时通过SQLEngineFactory.getSQLEngine(engineType)创建引擎并调用init完成 SQL 解析与校验每一行数据经transformRow交给sqlEngine.transformBySQL处理作业结束时close()释放引擎资源。行级错误分类Row-Level Error ClassificationSQLTransform实现了SupportRowLevelErrorClassifierSeaTunnelRow能够对单行处理失败进行分类见 SQLTransform.java表达式执行错误EXPRESSION_EXECUTE_ERROR→ 归类为ROW_ERROR行级错误可配置跳过等策略WHERE语句错误WHERE_STATEMENT_ERROR且不包含UNSUPPORTED_OPERATION原因 → 归类为ROW_ERROR其余错误 → 归类为SYSTEM_ERROR系统级错误。这为生产环境中的错误治理提供了依据例如当某一行数据因类型不合法导致表达式计算失败时作业可以据此决定是整作业失败还是按行级错误策略处理。Schema 变更处理对于 CDC 等动态 Schema 场景SQLTransform还实现了mapSchemaChangeEvent与setInputCatalogTable当上游发生ALTER TABLE如新增列时会置空缓存的 SQL 引擎迫使下一行数据到来时基于新 Schema 重新初始化避免因输出列数缓存过期导致ArrayIndexOutOfBoundsException源码注释对此有明确说明见 SQLTransform.java。对应的回归测试见 SQLMultiCatalogSchemaChangeTest.java 与 TransformChainLiveAlterTest.java。七、Sql 与 Calcite 变换的选型对比很多读者会困惑于SqlZeta 引擎与Calcite变换的区别。基于当前仓库两者对比如下维度Sql本文主题Calcite插件名SqlCalcite引擎Zeta内存 SQL 引擎基于 JSqlParserApache Calcite解析 → 校验 → 编译 → 执行入口类SQLTransform.javaCalciteTransform.java核心选项plugin_input/plugin_output/query/enginesql/table_transform/table_match_regex/row_error_handle_way行级错误处理通过SupportRowLevelErrorClassifier分类提供FAIL/SKIP/ROUTE_TO_TABLE三种策略多表 CDC单表变换支持table_transform按表覆盖 SQL、table_match_regex匹配多表向量类型视引擎函数支持而定内置向量 UDF如COSINE_DISTANCE、VECTOR_REDUCE向量类型内部映射为 VARBINARY两者都遵循逐行独立处理的原则JOIN、跨行聚合GROUP BY、SUM、COUNT均不支持。选择建议简单投影、过滤、函数计算用Sql即可满足需求需要多表 CDC 场景的按表 SQL 覆盖、向量运算或更标准化的 SQL 语义时可参考 Calcite Transform 文档。八、进阶资源与延伸阅读SQL Functions内置函数全清单本文query中可调用的全部字符串、数值、时间日期、系统与向量函数的语法、参数与示例SQL UDF自定义函数通过ZetaUDFSPI 机制为 Zeta 引擎注册自定义函数ZetaSQLEngine.loadUDFs使用ServiceLoader加载见 ZetaSQLEngine.javaCalcite Transform基于 Apache Calcite 的另一种 SQL 变换实现Transform 公共选项plugin_input/plugin_output等所有变换插件通用选项的说明Transforms Overview变换插件体系总览与新手推荐阅读顺序Multi Table Transform and Join Boundary了解多表变换与 JOIN 边界。源码与测试索引变换主类SQLTransform.java引擎工厂SQLEngineFactory.javaZeta 引擎实现ZetaSQLEngine.java函数求值实现ZetaSQLFunction.java、ZetaSQLType.java单元测试SQLTransformTest.java、ZetaSQLEngineTest.java、ZetaSQLFunctionTest.java九、常见问题速查Q1为什么我的 SQL 一提交就报Unsupported table join syntax/Unsupported GROUP BY syntax这是引擎的预期行为。Sql变换是逐行处理的内存引擎不支持 JOIN、GROUP BY、ORDER BY、LIMIT 等跨行或重排类语法请在作业启动前将这类逻辑拆解到 Source 端查询或改用其他方案如多表变换。Q2from dual是什么表名必须写dual吗不必。dual是引擎支持的占位表名参见 ZetaSQLEngine.java 中表名校验逻辑你也可以直接写from fakefake为plugin_input指定的表名或与上游 Catalog 表名一致。Q3Map 套 Map 的类型能不能取到内层 key 的值不能。引擎要求 Map 必须是查询路径上的最后一个结构mapstring, mapstring,string的内层 key 无法通过点号表达式访问。Q4engine配置成INTERNAL会怎样当前版本下ZETA与INTERNAL都会实例化ZetaSQLEngine行为一致见 SQLEngineFactory.java。该枚举为后续引擎扩展预留了空间。Q5上游 Schema 动态变化如 CDC 加列时会不会出问题不会。SQLTransform监听 Schema 变更事件并主动失效缓存引擎下一行数据会基于新 Schema 重新初始化执行计划见 SQLTransform.java相关行为有单元测试覆盖。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/18 9:31:35

DeepSeek Harness v0.5.2 插件加载失败根因与修复指南

1. 项目概述:这不是一个普通插件报错,而是本地大模型工作流的“心脏骤停”你点开 DeepSeek Harness 启动器,界面刚弹出来,底部状态栏突然飘出一行红字:“Plugin loading failed: Cannot resolve module ‘deepseek-har…

2026/9/18 9:31:35

uC/OS-II事件控制块、信号量与互斥量源码实战解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/18 10:36:56

代 Claude Docs 生成时,TaoToken 只提供 Key

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/18 10:36:56

STM32F103C8T6 入门实战:蓝牙遥控小车 GPIO/PWM/UART 速通教程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/18 10:36:56

RooCode 挂上 SumMCP.py 的 add / listdir,模型接口改填 TaoToken

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/9/18 10:31:56

DeepSeek V2到V3升级:兼容性排查与适配指南

简介:面向 DeepSeek 模型开发者与算法工程师的版本升级参考手册,聚焦从 V2 迁移到 V3 时的兼容性难题。内容按升级流程组织,先对比两代版本在技术架构、数据处理和模型输出上的差异,再依次说明硬件与软件环境准备、模型加载、分词…

2026/9/16 12:52:37

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/18 0:01:09

Google Colab 实战:运行模型、数据加载与报错排查

1. 为什么我劝你先搞懂 Colab 的运行模型1.1 Colab 到底是什么,跟本地跑代码差在哪Google Colab 简单说就是一台跑在浏览器里的 Linux 虚拟机,你打开一个 Notebook,背后就连上了一台带 GPU 的远程机器。你在单元格里敲的每一行 Python&#x…

2026/9/18 0:01:09

C语言数据类型与表达式详解

1. C语言数据与数据类型概述在C语言编程中,数据是程序处理的核心对象。理解数据的分类和特性是掌握C语言的基础。C语言中的数据主要分为四大类:常量、变量、表达式和函数。这些数据类型构成了C语言程序的基本元素,每种类型都有其独特的特性和…

2026/9/18 0:01:09

SQL时间字段指定时间段查询:区间语义、索引与时区避坑

上周排查一个线上问题&#xff0c;用户反馈"昨天的订单一条都没查到"&#xff0c;但数据库里明明躺着两千多条。最后定位下来&#xff0c;不是数据丢了&#xff0c;也不是接口挂了&#xff0c;而是那个查询条件把时间段写成了> 2024-05-20 00:00:00 AND < 2024…

2026/9/16 22:55:57

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行&#xff0c;Type-C接口算是典型的“看着简单&#xff0c;做起来全坑”的东西。光引脚就24个&#xff0c;高低速信号、电源、控制线全部塞在一个小小的连接器里&#xff0c;如果PCB布局不做规划&#xff0c;打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/16 22:56:09

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/16 22:56:16

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码