大数据流处理批处理数据工程【免费下载链接】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接口与函数解析的底层实现。什么是 ModulesModules 允许用户扩展 Flink 的内置对象built-in objects例如定义与 Flink 内置函数行为一致的自定义函数。它们是可插拔pluggable的Flink 提供了一些预置模块同时用户完全可以编写自己的模块。典型的使用场景包括用户定义自己的地理geo函数作为内置函数插入 Flink从而在 Flink SQL 和 Table API 中直接使用用户加载现成的 Hive 模块out-of-shelf把 Hive 的内置函数当作 Flink 的内置函数使用。更进一步一个模块还可以提供内置的 table source 和 sink 工厂规划部分从而禁用 Flink 基于 Java 服务提供者接口SPIService 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 表的DynamicTableSourceFactorygetTableSinkFactory()返回用于创建 sink 表的DynamicTableSinkFactory其中工厂方法的优先级规则源码注释明确给出持久化表对应 catalog 提供的工厂模块提供的工厂使用 Java SPI 发现的工厂。模块提供的工厂会按照模块加载顺序依次被调用第一个返回的工厂将被使用——这就是模块可以关闭 SPI 或影响临时表创建方式的底层依据。Module 类型CoreModuleCoreModule包含 Flink 的全部系统内置函数默认被加载且处于启用状态。其实现位于 CoreModule.java它持有BuiltInFunctionDefinitions.getDefinitions()返回的所有BuiltInFunctionDefinition在构造时以函数名大写化toUpperCase(Locale.ROOT)为 key 建立查找表listFunctions(false)会过滤掉内部函数isInternal()为 true 的getFunctionDefinition(name)则通过大写化后的名称进行不区分大小写的精确查找。CoreModule是单例CoreModule.INSTANCE这也解释了为什么它默认总是可用。HiveModuleHiveModule将 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-agggrouping、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(MapString, 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 精确实现内部用LinkedHashMapString, Module loadedModules保持加载顺序保证listFullModules()结果确定用ListString usedModules记录启用顺序构造时默认注册并启用 core 模块CoreModuleFactory.IDENTIFIERloadModule(name, module)同名模块已存在时抛出ValidationException否则加入usedModules与loadedModulesunloadModule(name)移除模块同时从启用列表移除不存在时抛异常useModules(names...)校验所有名字都已加载且不重复然后整体重置usedModules为新的声明顺序——这正是改变解析顺序/禁用未声明模块的实现机制getFunctionDefinition(name)按usedModules顺序遍历先检查listFunctions(true)中是否存在忽略大小写匹配的函数名命中即返回该模块的定义——这与文档中解析顺序决定同名函数归属的语义一一对应getFactory(selector)同样按启用顺序遍历模块返回第一个非空的工厂支撑模块工厂优先于 SPI 发现的机制。TableEnvironment的编程式入口在 TableEnvironmentImpl.javaloadModule、useModules、unloadModule三个方法均直接委托给ModuleManager。Namespace命名空间模块提供的对象被认为是 Flink 系统内置对象的一部分因此它们没有命名空间。这意味着模块导出的函数在 SQL 会话中是全局可见的不需要任何 schema 或 catalog 前缀。如何加载、卸载、使用和列出模块使用 SQL用户可以在 Table API 和 SQL CLI 中通过 SQL 来完成模块的加载/卸载/使用/列出操作。以下示例完整演示了从初始状态到加载 Hive 模块、调整解析顺序、禁用 core、卸载 Hive 的完整流程。JavaEnvironmentSettings 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 | // --------------------Scalaval 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 | // --------------------Pythonfrom 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 ValueCoreModulecoreHiveModulehivemodules: - 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。JavaEnvironmentSettings 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 | // --------------------Scalaval 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 | // --------------------Pythonfrom 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 配置中按nametype声明或通过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项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Bicep Registry Modules模块开发实战从零创建自定义模块Bicep Registry Modules模块开发实战从零创建自定义模块 概述 Bicep Registry Modules是微软Azure官方维护的BicIaC云原生sentence-transformers CrossEncoder 自定义模型开发指南模块链、保存加载机制与自定义模块实现sentence transformers CrossEncoder 自定义模型开发指南模块链、保存加载机制与自定义模块实现 本文基于 sentence tr人工智能NLPEmbedding微调机器学习Flink Table SQL LOAD 语句完全指南加载内置与自定义模块的原理与实战Flink Table SQL LOAD 语句完全指南加载内置与自定义模块的原理与实战 LOAD 语句 LOAD MODULE 是 Flink Table大数据流处理批处理数据工程上一篇Civitai OAuth 认证中枢切换上线后从配置核对到监控清理的完整运维 Checklist下一篇轻量级TTS方案选型espeak-ng与eSpeak、Flite的技术对比创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
