Flink Table Modules 模块机制详解:加载、解析顺序与自定义模块开发实战
2026/9/24 15:03:20 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/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

其中工厂方法的优先级规则(源码注释明确给出):

  1. 持久化表对应 catalog 提供的工厂;
  2. 模块提供的工厂;
  3. 使用 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_distrankrow_numberlagleadtumblehopsession系列等窗口/分析函数),这些函数不会被 Hive 模块覆盖,从而保证 Flink 自身窗口与时间属性语义不被破坏;
  • listFunctions()采用懒加载:首次调用时从hiveShim.listBuiltInFunctions()拉取全部内置函数,剔除黑名单,再补充groupingto_decimal等自定义实现;
  • 当配置项TABLE_EXEC_HIVE_NATIVE_AGG_FUNCTION_ENABLED开启时,sumcountavgminmax会被替换为 Flink 原生聚合函数实现(HiveSumAggFunction等),以支持 hash-agg;
  • groupinginternal_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,否则加入usedModulesloadedModules
  • unloadModule(name):移除模块,同时从启用列表移除,不存在时抛异常;
  • useModules(names...):校验所有名字都已加载且不重复,然后整体重置usedModules为新的声明顺序——这正是"改变解析顺序/禁用未声明模块"的实现机制;
  • getFunctionDefinition(name)usedModules顺序遍历,先检查listFunctions(true)中是否存在忽略大小写匹配的函数名,命中即返回该模块的定义——这与文档中"解析顺序决定同名函数归属"的语义一一对应;
  • getFactory(selector):同样按启用顺序遍历模块,返回第一个非空的工厂,支撑"模块工厂优先于 SPI 发现"的机制。

TableEnvironment的编程式入口在 TableEnvironmentImpl.java:loadModuleuseModulesunloadModule三个方法均直接委托给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属性来指明模块类型。开箱即支持以下类型:

ModuleType Value
CoreModulecore
HiveModulehive
modules: - name: core type: core - name: hive type: hive

⚠️ 注意:使用 SQL 时,模块名用于执行模块发现(module discovery),它会被解析为简单标识符(simple identifier)且区分大小写

使用 Java、Scala 或 Python 编程式管理

用户也可以通过编程方式(Java、Scala、Python)加载/卸载/使用/列出模块。与 SQL 方式对应的方法为loadModuleunloadModuleuseModuleslistModuleslistFullModules

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)

  1. usedModules的声明顺序遍历启用的模块;
  2. 对每个模块先调用listFunctions(true)(包含隐藏函数)做忽略大小写的匹配;
  3. 一旦命中,立即返回该模块的getFunctionDefinition(name),不再向后查找。

因此,USE MODULES hive, coreUSE MODULES core, hive会直接决定同名函数(例如 Hive 与 Flink 都提供的函数)的归属。这也解释了文档中的三条解析规则:启用状态 + 声明顺序共同决定解析结果;全部禁用则解析失败。

函数列举:SHOW MODULESSHOW 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())

因此,编写一个自定义模块通常需要两步:

  1. 实现Module接口,提供listFunctions()getFunctionDefinition(name)(以及可选地提供 source/sink 工厂);
  2. 实现ModuleFactory接口,声明factoryIdentifier()(作为 SQL 中的模块类型)、requiredOptions()/optionalOptions(),并实现createModule(Context)

之后即可在 SQL CLI 的 YAML 配置中按name+type声明,或通过LOAD MODULE语句动态加载。

注意事项与实践建议

  • 模块名区分大小写:SQL 中模块名被解析为简单标识符并区分大小写(如LOAD MODULE hiveLOAD 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

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载
上一篇:Civitai OAuth 认证中枢切换上线后:从配置核对到监控清理的完整运维 Checklist
下一篇:轻量级TTS方案选型:espeak-ng与eSpeak、Flite的技术对比

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询