SeaTunnel DynamicCompile 动态编译转换插件详解:运行时自定义数据处理与源码级原理

发布时间:2026/9/20 12:05:36

SeaTunnel DynamicCompile 动态编译转换插件详解:运行时自定义数据处理与源码级原理 SeaTunnel DynamicCompile 动态编译转换插件详解运行时自定义数据处理与源码级原理【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本指南围绕 SeaTunnel 的DynamicCompile转换插件展开讲解如何在作业运行时编译并执行 Groovy、Java、Scala 自定义代码实现灵活的行级数据处理如新增字段、改写字段、发起 RPC 请求、从外部数据源补全字段。读完本文你将掌握 DynamicCompile 的全部配置项、两种源码加载模式SOURCE_CODE / ABSOLUTE_PATH、必须实现的两个接口方法以及其底层基于 GroovyClassLoader、Janino、Scala REPL 的编译机制与 MD5 类缓存原理。插件简介与适用场景DynamicCompile是 SeaTunnel 提供的一种可编程转换插件允许用户在运行时编译并执行自定义代码来逐行处理数据。它把转换逻辑从配置文件提升到任意业务代码使得以下场景可以脱离内置插件直接实现自定义任何业务行为字段清洗、拼接、条件改写等基于现有行字段作为参数发起RPC 请求通过从其他数据源检索相关数据来扩展字段定义多个转换进行组合以便按业务边界拆分配置。需要特别说明的是插件官方文档对使用者给出安全提醒由于该插件会执行任意上传的代码必须确保服务安全防止攻击者上传破坏性代码同时如果转换逻辑过于复杂可能会影响整体性能需要在使用灵活性与吞吐量之间权衡。从源码结构看DynamicCompile 的实现位于 seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/dynamiccompile核心类DynamicCompileTransform继承自MultipleFieldOutputTransform本质上属于多字段输出型转换它既可以覆盖已有字段也可以向输出表追加新字段。配置属性总览DynamicCompile的全部配置项定义在 DynamicCompileTransformConfig.java 中属性如下nametyperequireddefault value说明source_codestringno内联源码当compile_pattern SOURCE_CODE时必填compile_languageEnumyes编译语言可选GROOVY/JAVA/SCALAcompile_patternEnumnoSOURCE_CODE源码提供方式可选SOURCE_CODE/ABSOLUTE_PATHabsolute_pathstringno服务器上 Java 或 Groovy 文件的绝对路径当compile_pattern ABSOLUTE_PATH时必填提示compile_language是唯一必填项compile_pattern缺省为SOURCE_CODE。通用选项 [string]plugin_input、plugin_output等转换插件的通用编排参数请参考 Transform 通用参数 了解详情。简单来说plugin_input声明当前 Transform 消费哪个上游数据集plugin_output把转换结果注册为命名数据集供下游引用。compile_language [Enum]可选值及编译实现GROOVY使用 Groovy 语法编写类底层通过GroovyClassLoader解析JAVA使用 Java 语法编写类底层通过 JaninoClassBodyEvaluator实现轻量级动态编译。注意 Janino 并非完整 JDK 编译器部分 Java 语法可能不受支持如泛型、注解、枚举等高级语法特性需避免使用SCALA目前支持 Zeta 引擎使用 Scala REPL 进行动态编译需要编写符合 Scala 语法的代码。对应的语言枚举定义在 CompileLanguage.java构造DynamicCompileTransform时会依据该枚举选择对应的解析器见 DynamicCompileTransform.java。compile_pattern [Enum]SOURCE_CODE从配置项的source_code属性直接读取内联源码ABSOLUTE_PATH从服务器上的文件路径absolute_path读取源码文件。选择SOURCE_CODE时source_code必填选择ABSOLUTE_PATH时absolute_path必填。这条互斥校验由 DynamicCompileTransformFactory.java 的OptionRule通过conditional条件规则实现并有对应的单元测试 DynamicCompileTransformFactoryTest.java 覆盖testValidSourceCodeConfig/testValidAbsolutePathConfig两种合法配置均可通过校验testMissingCompileLanguageFails缺少compile_language时抛出OptionValidationExceptiontestSourceCodeBlankFails/testAbsolutePathBlankFails对应源码/路径为空白字符串时校验失败testDefaultPatternMissingSourceCodeFails默认模式为SOURCE_CODE时未提供source_code同样校验失败。absolute_path [string]服务器上 Java 或 Groovy以及 Scala源码文件的绝对路径。插件运行时通过FileUtils.readFileToStr(Paths.get(absolute_path))读取文件内容之后走与内联源码完全相同的编译链路。注意路径是作业所在节点服务器上的路径需要保证该文件对执行进程可读。source_code [string]内联编写的源代码。这是 DynamicCompile 的核心必须实现两个方法Column[] getInlineOutputColumns(CatalogTable inputCatalogTable)Object[] getInlineOutputFieldValues(SeaTunnelRowAccessor inputRow)getInlineOutputColumns入参类型为CatalogTable返回Column[]。可以从入参的CatalogTable获取当前表的表结构。返回结果中如果字段已存在则按返回结果覆盖如果字段不存在则追加到现有表结构中。getInlineOutputFieldValues入参类型为SeaTunnelRowAccessor返回Object[]。可以从SeaTunnelRowAccessor获取当前行的数据如inputRow.getField(index)按索引取值执行自定义的数据处理逻辑。返回的数组长度必须与getInlineOutputColumns返回的列数一致且字段值顺序也要保持一致。从源码看这两个方法是通过ReflectionUtils.invoke反射调用的方法名常量即getInlineOutputColumns与getInlineOutputFieldValues每个输入行都会调用一次getInlineOutputFieldValues因此该方法中的逻辑应尽量轻量。第三方依赖加载如果自定义代码中引用了第三方依赖包请将它们放到${SEATUNNEL_HOME}/lib目录下如果使用 Spark 或 Flink 引擎则需要放到对应服务的 libs 目录下。添加依赖后必须重启集群服务才能重新加载。源码编译原理三种语言的三套解析器DynamicCompile 将源码 → Class的解析逻辑抽象为AbstractParse见 AbstractParse.java其唯一抽象方法parseClassSourceCode(String sourceCode)返回编译后的Class?。三种语言分别对应三个解析器解析器底层编译机制关键实现GroovyClassParser.javagroovy.lang.GroovyClassLoader静态单例 ClassLoaderparseClass(sourceCode)编译JavaClassParser.javaJaninoClassBodyEvaluatorcbe.cook(sourceCode)后getClazz()ScalaClassParser.javaScala REPLIMain静态初始化IMaincompileString后从 REPL 的 ClassLoader 加载类三个解析器都继承自 AbstractParser.java它提供了一层基于 MD5 的类缓存以源码内容的 MD5 摘要作为 key用ConcurrentHashMap缓存编译结果避免相同源码被反复编译。Scala 解析器还会通过正则(?:class|object)\s(\w)从源码中提取类名再通过 REPL 类加载器加载。此外DynamicCompileTransform在构造时还有一个兼容性检测如果源码中出现org.apache.seatunnel.transform.common.SeaTunnelRowAccessor旧的转型包路径则自动启用兼容模式将api.table.type包下的SeaTunnelRowAccessor包装为旧版访问器后传入保证历史代码无需修改即可运行见 DynamicCompileTransform.java 与getCompatibilityAccessor。示例为数据表新增字段并改写 age假设源端如FakeSource读取的数据表结构如下nameagecardJoy Ding20123May Ding20123Kin Dom30123Joy Dom30123我们将使用DynamicCompile对数据做两件事新增一列compile_language并将age20的行改写为age40。下面分别给出 Groovy、Java、源码文件路径与 Scala 四种写法。方式一Groovy 内联源码transform { DynamicCompile { plugin_input fake plugin_output groovy_out compile_languageGROOVY compile_patternSOURCE_CODE source_code import org.apache.seatunnel.api.table.catalog.Column import org.apache.seatunnel.api.table.type.SeaTunnelRowAccessor import org.apache.seatunnel.api.table.catalog.CatalogTable import org.apache.seatunnel.api.table.catalog.PhysicalColumn; import org.apache.seatunnel.api.table.type.*; import java.util.ArrayList; class demo { public Column[] getInlineOutputColumns(CatalogTable inputCatalogTable) { PhysicalColumn col1 PhysicalColumn.of( compile_language, BasicType.STRING_TYPE, 10L, true, , ); PhysicalColumn col2 PhysicalColumn.of( age, BasicType.INT_TYPE, 0L, false, false, ); return new Column[]{ col1, col2 }; } public Object[] getInlineOutputFieldValues(SeaTunnelRowAccessor inputRow) { Object[] fieldValues new Object[2]; // get age Object ageField inputRow.getField(1); fieldValues[0] GROOVY; if (Integer.parseInt(ageField.toString()) 20) { fieldValues[1] 40; } else { fieldValues[1] ageField; } return fieldValues; } }; } }方式二Java 内联源码transform { DynamicCompile { plugin_input fake plugin_output java_out compile_languageJAVA compile_patternSOURCE_CODE source_code import org.apache.seatunnel.api.table.catalog.Column; import org.apache.seatunnel.api.table.type.SeaTunnelRowAccessor; import org.apache.seatunnel.api.table.catalog.*; import org.apache.seatunnel.api.table.type.*; import java.util.ArrayList; public Column[] getInlineOutputColumns(CatalogTable inputCatalogTable) { PhysicalColumn col1 PhysicalColumn.of( compile_language, BasicType.STRING_TYPE, 10L, true, , ); PhysicalColumn col2 PhysicalColumn.of( age, BasicType.INT_TYPE, 0L, false, false, ); return new Column[]{ col1, col2 }; } public Object[] getInlineOutputFieldValues(SeaTunnelRowAccessor inputRow) { Object[] fieldValues new Object[2]; // get age Object ageField inputRow.getField(1); fieldValues[0] JAVA; if (Integer.parseInt(ageField.toString()) 20) { fieldValues[1] 40; } else { fieldValues[1] ageField; } return fieldValues; } } }方式三指定源码文件路径transform { DynamicCompile { plugin_input fake plugin_output groovy_out compile_languageGROOVY compile_patternABSOLUTE_PATH absolute_path/tmp/GroovyFile } }/tmp/GroovyFile中的类内容与方式一相同插件会读取该文件并执行同样的编译、调用流程。方式四Scala 内联源码transform { DynamicCompile { plugin_input fake plugin_output scala_out compile_languageSCALA compile_patternSOURCE_CODE source_code import org.apache.seatunnel.api.table.catalog.Column import org.apache.seatunnel.api.table.catalog.CatalogTable import org.apache.seatunnel.api.table.catalog.PhysicalColumn import org.apache.seatunnel.api.table.type.SeaTunnelRowAccessor import org.apache.seatunnel.api.table.type.BasicType import java.util.ArrayList class ScalaDemo { def getInlineOutputColumns(inputCatalogTable: CatalogTable): Array[Column] { val columns new ArrayList[Column]() val destColumn PhysicalColumn.of( compile_language, BasicType.STRING_TYPE, 10L, true, , ) columns.add(destColumn) columns.toArray(new ArrayColumn) } def getInlineOutputFieldValues(inputRow: SeaTunnelRowAccessor): Array[Object] { ArrayObject } } } }注意 Scala 示例中使用了反引号转义关键字org.apache.seatunnel.api.table.\type.SeaTunnelRowAccessor这是 Scala 中引用保留字type作为标识符的写法。同时 Scala 源码中需要定义class或object解析器会据此提取类名再交由 REPL 编译。转换结果执行 Groovy 配置后输出表groovy_out的数据将更新为nameagecardcompile_languageJoy Ding40123GROOVYMay Ding40123GROOVYKin Dom30123GROOVYJoy Dom30123GROOVY执行 Java 配置后输出表java_out的数据将更新为nameagecardcompile_languageJoy Ding40123JAVAMay Ding40123JAVAKin Dom30123JAVAJoy Dom30123JAVA可以看到compile_language是新增列getInlineOutputColumns返回结果中不存在于输入表的字段被追加age是已存在字段返回值按索引覆盖原值age20被改写为40。进阶实战在动态代码中调用 HTTP 接口DynamicCompile 的典型进阶用法是在getInlineOutputFieldValues中发起 RPC 调用将外部数据源返回值作为新字段输出。仓库中的端到端测试配置 single_dynamic_http_compile_transform.conf 演示了这一点class HttpDemo { public Column[] getInlineOutputColumns(CatalogTable inputCatalogTable) { ListColumn columns new ArrayList(); PhysicalColumn destColumn PhysicalColumn.of( DynamicCompile, BasicType.STRING_TYPE, 10, true, , ); columns.add(destColumn); return columns.toArray(new Column[0]); } public Object[] getInlineOutputFieldValues(SeaTunnelRowAccessor inputRow) { String body HttpUtil.get(http://mockserver:1080/v1/compile); Object[] fieldValues new Object[1]; fieldValues[0]body return fieldValues; } };该用例引入了cn.hutool.http.HttpUtil作为 HTTP 客户端依赖通过DependencyJar注入到测试容器的 Fake 插件 lib 下并配合 MockServer 容器模拟外部接口。这印证了文档中第三方依赖放入${SEATUNNEL_HOME}/lib后需重启的说明。多表支持与端到端验证多表Multi-Catalog支持DynamicCompile 通过 DynamicCompileMultiCatalogTransform.java 支持同时处理多张输入表对每张CatalogTable构建独立的DynamicCompileTransform并支持multi_tables、table_match_regex、rule_match_mode等通用多表选项见 DynamicCompileTransformFactory.java 的optionRule。端到端测试覆盖仓库在 TestDynamicCompileIT.java 中提供了完整的集成测试覆盖的配置文件位于 dynamic_compile/conf主要包括单语言测试single_dynamic_groovy_compile_transform.conf、single_dynamic_java_compile_transform.conf、single_dynamic_scala_compile_transform.conf路径模式测试single_groovy_path_compile.conf、single_java_path_compile.conf、single_scala_path_compile.conf多语言混合测试mixed_dynamic_groovy_java_compile_transform.conf、mixed_dynamic_groovy_scala_compile_transform.conf、mixed_dynamic_java_scala_compile_transform.conf、mixed_dynamic_all_compile_transform.conf多实例组合测试multiple_dynamic_java_compile_transform.conf、multiple_dynamic_groovy_compile_transform.conf、multiple_dynamic_scala_compile_transform.conf多表测试single_dynamic_java_compile_transform_multi_table.conf兼容模式测试single_dynamic_java_compile_transform_compatible.confHTTP 调用测试single_dynamic_http_compile_transform.conf。以 single_dynamic_java_compile_transform.conf 为例完整作业链路为FakeSource100 行字段id、name→DynamicCompile输出追加字符串列col1固定值test1→AssertSink 断言行数 ≥ 100 且col1字段非空、等于test1。这份配置可以直接作为编写自己作业时的最小可运行模板参考。使用注意事项与性能提示安全风险DynamicCompile 会执行任意的用户代码切勿在不可信环境下允许上传source_code并应通过权限控制限制absolute_path指向的文件语法限制Java 模式基于 Janino 编译部分 Java 语法不受支持复杂泛型、注解等特性建议改用 Groovy 或避免使用性能影响getInlineOutputFieldValues逐行调用逻辑应尽量轻量复杂的转换逻辑尤其同步 RPC 调用会拖慢吞吐建议控制动态代码的复杂度与调用频率字段一致性getInlineOutputFieldValues返回的数组长度与顺序必须严格对齐getInlineOutputColumns声明的列否则下游取数会错位依赖加载新增第三方 jar 到${SEATUNNEL_HOME}/lib或 Spark/Flink 的 libs后必须重启集群才能生效引擎支持SCALA语言目前仅支持 Zeta 引擎SeaTunnel 原生引擎使用其他引擎时请选用GROOVY或JAVA。小结DynamicCompile是 SeaTunnel 转换体系中灵活性最高的插件之一它把编码能力开放到配置层支持 Groovy / Java / Scala 三种语言的运行时编译提供内联源码与外部文件两种加载方式并通过反射调用两个约定方法实现改字段、加字段、查外部等任意行级业务逻辑。结合仓库源码可以看出其背后是 GroovyClassLoader、JaninoClassBodyEvaluator、Scala REPL 三套编译机制与 MD5 类缓存的组合配合工厂层的条件参数校验和多表Multi-Catalog支持构成了一个可用于生产的数据处理扩展点。使用时应始终牢记能力越强越要约束好代码来源与执行边界。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/20 12:05:36

Octop第一次对话指南:Web控制台Chat功能上手教程

Octop第一次对话指南:Web控制台Chat功能上手教程 【免费下载链接】Octop A smarter, self-hosted AI assistant — multi-user, multi-agent. 项目地址: https://gitcode.com/GitHub_Trending/oct/Octop Octop 是一个开源、自托管的多用户多 Agent AI 助手&a…

2026/9/20 13:05:40

高项论文写作:如何避免雷区并打造差异化内容

1. 高项论文写作的现状与痛点每次看到考生们拿着千篇一律的论文模板来咨询修改意见,我都忍不住想提醒:评审专家每年要看上千份论文,那些老掉牙的案例和套路化的表达,早就让他们审美疲劳了。去年有位考生用了某培训机构的"万能…

2026/9/20 13:05:40

大学邮箱第三方客户端配置与安全指南

1. 大学邮箱第三方客户端配置全指南作为使用大学邮箱多年的老用户,我深知通过手机客户端实时查收学校通知和学术邮件的必要性。但很多同学在配置第三方客户端时总会遇到各种问题,今天我就把完整的配置流程和避坑要点整理出来。大学邮箱系统通常基于Corem…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/20 4:54:47

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

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

2026/9/20 5:01:23

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

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

2026/9/20 5:09:33

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

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

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

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

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