- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
导读
Flink Table Modules(模块)是 Table API / SQL 中用于扩展系统内置对象(如内置函数)的可插拔机制:它允许你把自定义函数注册为"类内置函数",在 SQL 和 Table API 中直接使用,也可以一键挂载 Hive 模块来复用 Hive 的内置函数生态。读完本文,你将掌握 Module 的核心概念与三种模块类型、模块的生命周期(加载/启用/禁用/卸载)与同名函数解析顺序规则,并能通过 SQL、Table API 或 SQL Client 的 YAML 配置完整操控模块,最后深入源码理解Module/ModuleFactory接口与函数解析的底层实现。
什么是 Modules
Modules 允许用户扩展 Flink 的内置对象(built-in objects),例如定义与 Flink 内置函数行为一致的自定义函数。它们是可插拔(pluggable)的:Flink 提供了一些预置模块,同时用户完全可以编写自己的模块。
典型的使用场景包括:
- 用户定义自己的地理(geo)函数,作为内置函数插入 Flink,从而在 Flink SQL 和 Table API 中直接使用;
- 用户加载现成的 Hive 模块(out-of-shelf),把 Hive 的内置函数当作 Flink 的内置函数使用。
更进一步,一个模块还可以提供内置的 table source 和 sink 工厂(规划部分),从而禁用 Flink 基于 Java 服务提供者接口(SPI,Service Provider Interfaces)的默认发现机制,或者在没有对应 catalog 的情况下影响临时表连接器(connector)的创建方式。
从源码上看,模块定义在 Module.java 中:其 Javadoc 明确说明"Modules define a set of metadata, including functions, user defined types, operators, rules, etc. Metadata from modules are regarded as built-in or system metadata"。Module接口标注为@PublicEvolving,核心方法包括:
| 方法 | 作用 |
|---|---|
listFunctions()/listFunctions(boolean includeHiddenFunctions) | 列出模块中所有函数名,默认排除内部/隐藏函数 |
getFunctionDefinition(String name) | 按名称返回可选的FunctionDefinition(包含隐藏函数) |
getTableSourceFactory() | 返回用于创建 source 表的DynamicTableSourceFactory |
getTableSinkFactory() | 返回用于创建 sink 表的DynamicTableSinkFactory |
其中工厂方法的优先级规则(源码注释明确给出):
- 持久化表对应 catalog 提供的工厂;
- 模块提供的工厂;
- 使用 Java SPI 发现的工厂。
模块提供的工厂会按照模块加载顺序依次被调用,第一个返回的工厂将被使用——这就是模块可以"关闭 SPI 或影响临时表创建方式"的底层依据。
Module 类型
CoreModule
CoreModule包含 Flink 的全部系统(内置)函数,默认被加载且处于启用状态。其实现位于 CoreModule.java:它持有BuiltInFunctionDefinitions.getDefinitions()返回的所有BuiltInFunctionDefinition,在构造时以函数名(大写化,toUpperCase(Locale.ROOT))为 key 建立查找表;listFunctions(false)会过滤掉内部函数(isInternal()为 true 的),getFunctionDefinition(name)则通过大写化后的名称进行不区分大小写的精确查找。CoreModule是单例(CoreModule.INSTANCE),这也解释了为什么它默认总是可用。
HiveModule
HiveModule将 Hive 的内置函数作为 Flink 的系统函数提供给 SQL 和 Table API 用户。Flink 的 Hive 文档 提供了搭建该模块的完整细节。
其实现位于 HiveModule.java,几个值得注意的源码细节:
- 构造函数要求非空
hiveVersion,并通过HiveShimLoader.loadHiveShim(hiveVersion)按版本加载对应的 Hive shim,实现多 Hive 版本兼容; - 维护了一个
BUILT_IN_FUNC_BLACKLIST黑名单(如cume_dist、rank、row_number、lag、lead、tumble、hop、session系列等窗口/分析函数),这些函数不会被 Hive 模块覆盖,从而保证 Flink 自身窗口与时间属性语义不被破坏; listFunctions()采用懒加载:首次调用时从hiveShim.listBuiltInFunctions()拉取全部内置函数,剔除黑名单,再补充grouping、to_decimal等自定义实现;- 当配置项
TABLE_EXEC_HIVE_NATIVE_AGG_FUNCTION_ENABLED开启时,sum、count、avg、min、max会被替换为 Flink 原生聚合函数实现(HiveSumAggFunction等),以支持 hash-agg; grouping、internal_interval等函数被覆盖为 Flink 侧的兼容实现。
HiveModule的工厂 HiveModuleFactory.java 的标识符(factoryIdentifier())为hive,其可选配置项定义在 HiveModuleOptions.java 中,即hive-version选项(字符串类型,无默认值);未显式指定时回退到HiveShimLoader.getHiveVersion()。
用户自定义模块(User-Defined Module)
用户可以通过实现Module接口开发自定义模块。要在 SQL CLI 中使用自定义模块,需要同时开发模块本身以及对应的模块工厂——即实现ModuleFactory接口。
ModuleFactory定义在 ModuleFactory.java,标注为@PublicEvolving。一个模块工厂定义了一组属性(properties),用于在 SQL CLI 启动(bootstrap)时配置该模块。这些属性被传递给一个发现服务(discovery service),服务尝试将属性与某个ModuleFactory匹配,并实例化对应的模块实例。
ModuleFactory的关键要素:
createModule(Context context):根据上下文创建并配置模块;Context接口提供getOptions()(创建模块的选项,实现方应校验)、getConfiguration()(当前会话的只读配置)、getClassLoader()(当前会话的类加载器,可用于发现嵌套工厂);factoryIdentifier():工厂唯一标识符,SQL 中的模块类型即与之对应;requiredOptions()/optionalOptions():必选/可选配置项集合,用于发现与校验。
老旧的createModule(Map<String, String>)(基于TableFactory栈)已被标记@Deprecated,新实现应基于Factory栈。
Module 生命周期与解析顺序
一个模块可以被加载(load)、启用(use/enable)、禁用(disable)和卸载(unload)。
- 当
TableEnvironment初次加载一个模块时,默认会启用它; - Flink 支持多个模块并存,并跟踪加载顺序来解析元数据;
- Flink 只会在启用的模块中解析函数。
当两个模块中存在同名函数时,有三种情况:
- 两个模块都被启用:Flink 按模块的解析顺序(resolution order)解析函数;
- 其中一个被禁用:Flink 解析到被启用的那个模块;
- 两个都被禁用:Flink 无法解析该函数。
用户可以通过不同的声明顺序改变解析顺序。例如,通过USE MODULES hive, core让 Flink 优先在 Hive 中查找函数。
此外,用户也可以"不声明"某个模块来禁用它。例如USE MODULES hive会禁用 core 模块(强烈不推荐禁用 core 模块)。注意:禁用模块并不会卸载它,用户可以通过再次使用它来重新启用,例如USE MODULES core, hive会把 core 模块带回并置于首位。一个模块只有在其已加载(loaded)的状态下才能被启用;使用一个未加载的模块会抛出异常。最后,用户可以卸载(unload)一个模块。
禁用与卸载的区别在于:TableEnvironment仍然保留被禁用的模块,用户可以通过列出所有已加载的模块来查看被禁用的模块。
在源码层面,这些行为由 ModuleManager.java 精确实现:
- 内部用
LinkedHashMap<String, Module> loadedModules保持加载顺序(保证listFullModules()结果确定),用List<String> usedModules记录启用顺序; - 构造时默认注册并启用 core 模块(
CoreModuleFactory.IDENTIFIER); loadModule(name, module):同名模块已存在时抛出ValidationException,否则加入usedModules与loadedModules;unloadModule(name):移除模块,同时从启用列表移除,不存在时抛异常;useModules(names...):校验所有名字都已加载且不重复,然后整体重置usedModules为新的声明顺序——这正是"改变解析顺序/禁用未声明模块"的实现机制;getFunctionDefinition(name):按usedModules顺序遍历,先检查listFunctions(true)中是否存在忽略大小写匹配的函数名,命中即返回该模块的定义——这与文档中"解析顺序决定同名函数归属"的语义一一对应;getFactory(selector):同样按启用顺序遍历模块,返回第一个非空的工厂,支撑"模块工厂优先于 SPI 发现"的机制。
TableEnvironment的编程式入口在 TableEnvironmentImpl.java:loadModule、useModules、unloadModule三个方法均直接委托给ModuleManager。
Namespace(命名空间)
模块提供的对象被认为是 Flink 系统(内置)对象的一部分,因此它们没有命名空间。这意味着模块导出的函数在 SQL 会话中是全局可见的,不需要任何 schema 或 catalog 前缀。
如何加载、卸载、使用和列出模块
使用 SQL
用户可以在 Table API 和 SQL CLI 中通过 SQL 来完成模块的加载/卸载/使用/列出操作。以下示例完整演示了从初始状态到加载 Hive 模块、调整解析顺序、禁用 core、卸载 Hive 的完整流程。
Java
EnvironmentSettings settings = EnvironmentSettings.inStreamingMode(); TableEnvironment tableEnv = TableEnvironment.create(settings); // Show initially loaded and enabled modules tableEnv.executeSql("SHOW MODULES").print(); // +-------------+ // | module name | // +-------------+ // | core | // +-------------+ tableEnv.executeSql("SHOW FULL MODULES").print(); // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // +-------------+------+ // Load a hive module tableEnv.executeSql("LOAD MODULE hive WITH ('hive-version' = '...')"); // Show all enabled modules tableEnv.executeSql("SHOW MODULES").print(); // +-------------+ // | module name | // +-------------+ // | core | // | hive | // +-------------+ // Show all loaded modules with both name and use status tableEnv.executeSql("SHOW FULL MODULES").print(); // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // | hive | true | // +-------------+------+ // Change resolution order tableEnv.executeSql("USE MODULES hive, core"); tableEnv.executeSql("SHOW MODULES").print(); // +-------------+ // | module name | // +-------------+ // | hive | // | core | // +-------------+ tableEnv.executeSql("SHOW FULL MODULES").print(); // +-------------+------+ // | module name | used | // +-------------+------+ // | hive | true | // | core | true | // +-------------+------+ // Disable core module tableEnv.executeSql("USE MODULES hive"); tableEnv.executeSql("SHOW MODULES").print(); // +-------------+ // | module name | // +-------------+ // | hive | // +-------------+ tableEnv.executeSql("SHOW FULL MODULES").print(); // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | true | // | core | false | // +-------------+-------+ // Unload hive module tableEnv.executeSql("UNLOAD MODULE hive"); tableEnv.executeSql("SHOW MODULES").print(); // Empty set tableEnv.executeSql("SHOW FULL MODULES").print(); // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | false | // +-------------+-------+Scala
val settings = EnvironmentSettings.inStreamingMode() val tableEnv = TableEnvironment.create(setting) // Show initially loaded and enabled modules tableEnv.executeSql("SHOW MODULES").print() // +-------------+ // | module name | // +-------------+ // | core | // +-------------+ tableEnv.executeSql("SHOW FULL MODULES").print() // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // +-------------+------+ // Load a hive module tableEnv.executeSql("LOAD MODULE hive WITH ('hive-version' = '...')") // Show all enabled modules tableEnv.executeSql("SHOW MODULES").print() // +-------------+ // | module name | // +-------------+ // | core | // | hive | // +-------------+ // Show all loaded modules with both name and use status tableEnv.executeSql("SHOW FULL MODULES") // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // | hive | true | // +-------------+------+ // Change resolution order tableEnv.executeSql("USE MODULES hive, core") tableEnv.executeSql("SHOW MODULES").print() // +-------------+ // | module name | // +-------------+ // | hive | // | core | // +-------------+ tableEnv.executeSql("SHOW FULL MODULES").print() // +-------------+------+ // | module name | used | // +-------------+------+ // | hive | true | // | core | true | // +-------------+------+ // Disable core module tableEnv.executeSql("USE MODULES hive") tableEnv.executeSql("SHOW MODULES").print() // +-------------+ // | module name | // +-------------+ // | hive | // +-------------+ tableEnv.executeSql("SHOW FULL MODULES").print() // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | true | // | core | false | // +-------------+-------+ // Unload hive module tableEnv.executeSql("UNLOAD MODULE hive") tableEnv.executeSql("SHOW MODULES").print() // Empty set tableEnv.executeSql("SHOW FULL MODULES").print() // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | false | // +-------------+-------+Python
from pyflink.table import * # environment configuration settings = EnvironmentSettings.inStreamingMode() t_env = TableEnvironment.create(settings) # Show initially loaded and enabled modules t_env.execute_sql("SHOW MODULES").print() # +-------------+ # | module name | # +-------------+ # | core | # +-------------+ t_env.execute_sql("SHOW FULL MODULES").print() # +-------------+------+ # | module name | used | # +-------------+------+ # | core | true | # +-------------+------+ # Load a hive module t_env.execute_sql("LOAD MODULE hive WITH ('hive-version' = '...')") # Show all enabled modules t_env.execute_sql("SHOW MODULES").print() # +-------------+ # | module name | # +-------------+ # | core | # | hive | # +-------------+ # Show all loaded modules with both name and use status t_env.execute_sql("SHOW FULL MODULES").print() # +-------------+------+ # | module name | used | # +-------------+------+ # | core | true | # | hive | true | # +-------------+------+ # Change resolution order t_env.execute_sql("USE MODULES hive, core") t_env.execute_sql("SHOW MODULES").print() # +-------------+ # | module name | # +-------------+ # | hive | # | core | # +-------------+ t_env.execute_sql("SHOW FULL MODULES").print() # +-------------+------+ # | module name | used | # +-------------+------+ # | hive | true | # | core | true | # +-------------+------+ # Disable core module t_env.execute_sql("USE MODULES hive") t_env.execute_sql("SHOW MODULES").print() # +-------------+ # | module name | # +-------------+ # | hive | # +-------------+ t_env.execute_sql("SHOW FULL MODULES").print() # +-------------+-------+ # | module name | used | # +-------------+-------+ # | hive | true | # | core | false | # +-------------+-------+ # Unload hive module t_env.execute_sql("UNLOAD MODULE hive") t_env.execute_sql("SHOW MODULES").print() # Empty set t_env.execute_sql("SHOW FULL MODULES").print() # +-------------+-------+ # | module name | used | # +-------------+-------+ # | hive | false | # +-------------+-------+SQL Client
-- Show initially loaded and enabled modules Flink SQL> SHOW MODULES; +-------------+ | module name | +-------------+ | core | +-------------+ 1 row in set Flink SQL> SHOW FULL MODULES; +-------------+------+ | module name | used | +-------------+------+ | core | true | +-------------+------+ 1 row in set -- Load a hive module Flink SQL> LOAD MODULE hive WITH ('hive-version' = '...'); -- Show all enabled modules Flink SQL> SHOW MODULES; +-------------+ | module name | +-------------+ | core | | hive | +-------------+ 2 rows in set -- Show all loaded modules with both name and use status Flink SQL> SHOW FULL MODULES; +-------------+------+ | module name | used | +-------------+------+ | core | true | | hive | true | +-------------+------+ 2 rows in set -- Change resolution order Flink SQL> USE MODULES hive, core ; Flink SQL> SHOW MODULES; +-------------+ | module name | +-------------+ | hive | | core | +-------------+ 2 rows in set Flink SQL> SHOW FULL MODULES; +-------------+------+ | module name | used | +-------------+------+ | hive | true | | core | true | +-------------+------+ 2 rows in set -- Unload hive module Flink SQL> UNLOAD MODULE hive; Flink SQL> SHOW MODULES; Empty set Flink SQL> SHOW FULL MODULES; +-------------+-------+ | module name | used | +-------------+-------+ | hive | false | +-------------+-------+ 1 row in set通过 YAML 配置 SQL Client
在 SQL Client 的 YAML 配置文件中定义的所有模块都必须提供type属性来指明模块类型。开箱即支持以下类型:
| Module | Type Value |
|---|---|
| CoreModule | core |
| HiveModule | hive |
modules: - name: core type: core - name: hive type: hive⚠️ 注意:使用 SQL 时,模块名用于执行模块发现(module discovery),它会被解析为简单标识符(simple identifier)且区分大小写。
使用 Java、Scala 或 Python 编程式管理
用户也可以通过编程方式(Java、Scala、Python)加载/卸载/使用/列出模块。与 SQL 方式对应的方法为loadModule、unloadModule、useModules、listModules、listFullModules。
Java
EnvironmentSettings settings = EnvironmentSettings.inStreamingMode(); TableEnvironment tableEnv = TableEnvironment.create(settings); // Show initially loaded and enabled modules tableEnv.listModules(); // +-------------+ // | module name | // +-------------+ // | core | // +-------------+ tableEnv.listFullModules(); // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // +-------------+------+ // Load a hive module tableEnv.loadModule("hive", new HiveModule()); // Show all enabled modules tableEnv.listModules(); // +-------------+ // | module name | // +-------------+ // | core | // | hive | // +-------------+ // Show all loaded modules with both name and use status tableEnv.listFullModules(); // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // | hive | true | // +-------------+------+ // Change resolution order tableEnv.useModules("hive", "core"); tableEnv.listModules(); // +-------------+ // | module name | // +-------------+ // | hive | // | core | // +-------------+ tableEnv.listFullModules(); // +-------------+------+ // | module name | used | // +-------------+------+ // | hive | true | // | core | true | // +-------------+------+ // Disable core module tableEnv.useModules("hive"); tableEnv.listModules(); // +-------------+ // | module name | // +-------------+ // | hive | // +-------------+ tableEnv.listFullModules(); // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | true | // | core | false | // +-------------+-------+ // Unload hive module tableEnv.unloadModule("hive"); tableEnv.listModules(); // Empty set tableEnv.listFullModules(); // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | false | // +-------------+-------+Scala
val settings = EnvironmentSettings.inStreamingMode() val tableEnv = TableEnvironment.create(setting) // Show initially loaded and enabled modules tableEnv.listModules() // +-------------+ // | module name | // +-------------+ // | core | // +-------------+ tableEnv.listFullModules() // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // +-------------+------+ // Load a hive module tableEnv.loadModule("hive", new HiveModule()) // Show all enabled modules tableEnv.listModules() // +-------------+ // | module name | // +-------------+ // | core | // | hive | // +-------------+ // Show all loaded modules with both name and use status tableEnv.listFullModules() // +-------------+------+ // | module name | used | // +-------------+------+ // | core | true | // | hive | true | // +-------------+------+ // Change resolution order tableEnv.useModules("hive", "core") tableEnv.listModules() // +-------------+ // | module name | // +-------------+ // | hive | // | core | // +-------------+ tableEnv.listFullModules() // +-------------+------+ // | module name | used | // +-------------+------+ // | hive | true | // | core | true | // +-------------+------+ // Disable core module tableEnv.useModules("hive") tableEnv.listModules() // +-------------+ // | module name | // +-------------+ // | hive | // +-------------+ tableEnv.listFullModules() // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | true | // | core | false | // +-------------+-------+ // Unload hive module tableEnv.unloadModule("hive") tableEnv.listModules() // Empty set tableEnv.listFullModules() // +-------------+-------+ // | module name | used | // +-------------+-------+ // | hive | false | // +-------------+-------+Python
from pyflink.table import * # environment configuration settings = EnvironmentSettings.inStreamingMode() t_env = TableEnvironment.create(settings) # Show initially loaded and enabled modules t_env.list_modules() # +-------------+ # | module name | # +-------------+ # | core | # +-------------+ t_env.list_full_modules() # +-------------+------+ # | module name | used | # +-------------+------+ # | core | true | # +-------------+------+ # Load a hive module t_env.load_module("hive", HiveModule()) # Show all enabled modules t_env.list_modules() # +-------------+ # | module name | # +-------------+ # | core | # | hive | # +-------------+ # Show all loaded modules with both name and use status t_env.list_full_modules() # +-------------+------+ # | module name | used | # +-------------+------+ # | core | true | # | hive | true | # +-------------+------+ # Change resolution order t_env.use_modules("hive", "core") t_env.list_modules() # +-------------+ # | module name | # +-------------+ # | hive | # | core | # +-------------+ t_env.list_full_modules() # +-------------+------+ # | module name | used | # +-------------+------+ # | hive | true | # | core | true | # +-------------+------+ # Disable core module t_env.use_modules("hive") t_env.list_modules() # +-------------+ # | module name | # +-------------+ # | hive | # +-------------+ t_env.list_full_modules() # +-------------+-------+ # | module name | used | # +-------------+-------+ # | hive | true | # | core | false | # +-------------+-------+ # Unload hive module t_env.unload_module("hive") t_env.list_modules() # Empty set t_env.list_full_modules() # +-------------+-------+ # | module name | used | # +-------------+-------+ # | hive | false | # +-------------+-------+深入源码:模块函数的解析与工厂发现原理
函数解析:按启用顺序"先到先得"
结合前面 ModuleManager.java 的实现可以看到,一次 SQL 中的函数调用最终会走到ModuleManager.getFunctionDefinition(name):
- 按
usedModules的声明顺序遍历启用的模块; - 对每个模块先调用
listFunctions(true)(包含隐藏函数)做忽略大小写的匹配; - 一旦命中,立即返回该模块的
getFunctionDefinition(name),不再向后查找。
因此,USE MODULES hive, core与USE MODULES core, hive会直接决定同名函数(例如 Hive 与 Flink 都提供的函数)的归属。这也解释了文档中的三条解析规则:启用状态 + 声明顺序共同决定解析结果;全部禁用则解析失败。
函数列举:SHOW MODULES与SHOW FULL MODULES的差异
ModuleManager.listModules()仅返回usedModules(启用中的模块,按解析顺序);listFullModules()则返回全部已加载模块的ModuleEntry(模块名 + 是否启用),启用中的在前(按解析顺序)、禁用的在后。这正好对应 SQL 中SHOW MODULES只显示启用的模块、而SHOW FULL MODULES额外带出used状态列的行为——也是"禁用不等于卸载"这一语义的可观测体现。
模块工厂发现:自定义模块接入 SQL CLI 的路径
当在 SQL Client / YAML 中声明type: hive或执行LOAD MODULE hive WITH (...)时,SQL 层会把模块名(如hive)作为标识符交给发现服务,服务将其与ModuleFactory.factoryIdentifier()进行匹配,找到后调用createModule(Context)完成实例化。以HiveModuleFactory为例,它声明factoryIdentifier() = "hive"、可选配置项hive-version,并在createModule中通过FactoryUtil.createModuleFactoryHelper做选项校验,最终new HiveModule(hiveVersion, context.getConfiguration(), context.getClassLoader())。
因此,编写一个自定义模块通常需要两步:
- 实现
Module接口,提供listFunctions()与getFunctionDefinition(name)(以及可选地提供 source/sink 工厂); - 实现
ModuleFactory接口,声明factoryIdentifier()(作为 SQL 中的模块类型)、requiredOptions()/optionalOptions(),并实现createModule(Context)。
之后即可在 SQL CLI 的 YAML 配置中按name+type声明,或通过LOAD MODULE语句动态加载。
注意事项与实践建议
- 模块名区分大小写:SQL 中模块名被解析为简单标识符并区分大小写(如
LOAD MODULE hive与LOAD MODULE Hive是不同处理路径); - 禁用不等于卸载:
USE MODULES hive只是把 core 从启用列表移除,core 仍处于已加载状态,可通过USE MODULES core, hive重新启用;而UNLOAD MODULE才会真正移除模块; - 不要轻易禁用 core 模块:文档明确"强烈不推荐禁用 core 模块",因为 Flink 的系统内置函数都来自
CoreModule,禁用后大量内置函数将不可解析; - 同名函数的行为差异:Hive 模块自带黑名单,不会覆盖 Flink 的窗口函数与分析函数;若你的自定义模块与既有模块函数重名,务必通过
USE MODULES明确解析顺序,避免语义混乱; - 模块的工厂能力:如果希望完全掌控临时表 source/sink 的创建,可以在
Module中实现getTableSourceFactory()/getTableSinkFactory(),按加载顺序优先于 Java SPI 生效(工厂优先级为:catalog > module > SPI)。
以上所有 SQL 语句与编程 API 在 Table API、SQL CLI 中均可直接运行验证;模块机制的完整源码可继续阅读 Module.java、CoreModule.java、HiveModule.java、ModuleManager.java 与 TableEnvironmentImpl.java。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Bicep Registry Modules模块开发实战:从零创建自定义模块
Bicep Registry Modules模块开发实战:从零创建自定义模块 概述 Bicep Registry Modules是微软Azure官方维护的Bic
IaC云原生sentence-transformers CrossEncoder 自定义模型开发指南:模块链、保存加载机制与自定义模块实现
sentence transformers CrossEncoder 自定义模型开发指南:模块链、保存加载机制与自定义模块实现 本文基于 sentence tr
人工智能NLPEmbedding微调机器学习Flink Table SQL LOAD 语句完全指南:加载内置与自定义模块的原理与实战
Flink Table SQL LOAD 语句完全指南:加载内置与自定义模块的原理与实战 LOAD 语句( LOAD MODULE )是 Flink Table
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考