ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Flink 集成 Hive 函数:HiveModule、原生聚合函数与 Hive UDF 完整使用指南

Apache Flink 集成 Hive 函数:HiveModule、原生聚合函数与 Hive UDF 完整使用指南 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载本指南系统讲解如何在 Flink SQL 与 Table API 中复用 Hive 的函数生态通过HiveModule将 Hive 内置函数加载为 Flink 系统内置函数通过table.exec.hive.native-agg-function.enabled启用基于 hash 聚合的原生 Hive 聚合函数以获得性能提升以及在 Flink 中直接调用 Hive 用户自定义函数UDF/GenericUDF/GenericUDTF/UDAF。读完本文你将掌握 Hive 函数在 Flink 中的三种接入方式、相关配置参数及其底层实现原理。Hive 函数集成的三种途径在 Flink 中集成 Hive 函数主要有三条路径本文分别展开Hive 内置函数Built-in Functions通过加载HiveModule把 Hive 内置函数作为 Flink 系统级内置函数提供给 Flink SQL 和 Table API 用户使用原生 Hive 聚合函数Native Hive Aggregate FunctionsFlink 1.17 起引入的一组由 Flink 原生实现的 Hive 聚合函数sum/count/avg/min/max可借助 hash 聚合算子执行显著提升聚合性能Hive 用户自定义函数UDF复用 Hive Metastore 中注册的 UDF、GenericUDF、GenericUDTF、UDAF、GenericUDAFResolver2在 Flink SQL 中直接调用。后文将结合 flink-connector-hive 模块的源码说明每种方式的实现细节。通过 HiveModule 使用 Hive 内置函数HiveModule将 Hive 的内置函数以 Flink 系统内置函数的形式提供给 Flink SQL 与 Table API 用户。加载该模块后Flink 会话中即可直接调用 Hive 内置函数无需额外注册。详细的模块机制可参考 Flink Modules 文档。加载 HiveModule以下分别给出 Java、Scala、Python 与 SQL Client 四种加载方式。JavaString name myhive; String version 2.3.4; tableEnv.loadModule(name, new HiveModule(version));Scalaval name myhive val version 2.3.4 tableEnv.loadModule(name, new HiveModule(version));Pythonfrom pyflink.table.module import HiveModule name myhive version 2.3.4 t_env.load_module(name, HiveModule(version))SQL ClientLOAD MODULE hive WITH (hive-version 2.3.4);其中version/hive-version为待加载的 Hive 版本号例如2.3.4。从源码看HiveModule的构造方法会基于该版本号通过HiveShimLoader.loadHiveShim(hiveVersion)加载对应版本的 Hive Shim见 HiveModule.java因此版本号必须与 classpath 中的 Hive 依赖匹配。模块解析顺序让 Hive 函数优先于 Core 函数Flink 默认加载并启用CoreModule包含 Flink 全部系统内置函数。当同时加载了HiveModule时Flink 会按照模块的声明顺序解析同名函数先声明者优先。若希望 Hive 内置函数优先于 Flink 核心函数被解析可执行USE MODULES hive, core;该语句将hive模块置于core之前此后同名函数如sum、count将优先解析到 Hive 实现。也可以只启用hive模块USE MODULES hive但文档不建议禁用 core 模块。更多模块生命周期load/enable/disable/unload与解析规则参见 Modules 文档。源码视角HiveModule 如何暴露函数从 HiveModule.java 可以看到其核心逻辑listFunctions()懒加载 Hive 内置函数名集合通过hiveShim.listBuiltInFunctions()并移除黑名单函数见 HiveModule.java。黑名单中的函数如rank、row_number、lag、tumble等窗口/分析函数不会被 HiveModule 覆盖从而保证 Flink 自身的窗口与解析函数行为不被破坏getFunctionDefinition(name)对非黑名单函数通过HiveFunctionDefinitionFactory把 Hive 函数包装为 Flink 的FunctionDefinition并额外覆盖了grouping、internal_interval、to_decimal等特殊函数见 HiveModule.java。注意事项旧版本 Hive 内置函数的线程安全问题部分旧版本 Hive 内置函数存在线程安全问题HIVE-16183。官方建议用户自行对 Hive 打补丁修复。在使用前请确认所依赖的 Hive 版本中相关内置函数无此类问题。使用原生 Hive 聚合函数Native Hive Aggregate Functions当HiveModule的优先级高于CoreModule时Flink 会优先尝试使用 Hive 内置函数。对于 Hive 内置聚合函数Flink 此前只能使用基于排序的聚合算子sort-based aggregation。从 Flink 1.17 开始Flink 引入了一组原生 Hive 聚合函数可以使用基于哈希的聚合算子hash-based aggregation执行从而带来显著的聚合性能提升。支持的函数与配置项当前仅支持以下五个函数未来会支持更多sumcountavgminmax通过开启以下作业级配置项启用原生聚合函数KeyDefaultTypeDescriptiontable.exec.hive.native-agg-function.enabledfalseBoolean启用原生聚合函数可使用 hash 聚合策略提升聚合性能。这是一个作业级job-level配置。在 SQL Client 中开启方式示例SET table.exec.hive.native-agg-function.enabled true;注意Attention当前原生聚合函数的能力与 Hive 内置聚合函数并非完全对齐例如部分数据类型尚不支持。如果性能不是瓶颈无需开启该选项。注意Attention通过 SqlClient 使用时table.exec.hive.native-agg-function.enabled目前不能按作业per job单独开启仅支持模块级module level配置。用户应当先开启该选项再加载 HiveModule。该限制将在未来版本修复。源码视角原生聚合函数的实现配置项定义于 HiveOptions.javapublic static final ConfigOptionBoolean TABLE_EXEC_HIVE_NATIVE_AGG_FUNCTION_ENABLED key(table.exec.hive.native-agg-function.enabled) .booleanType() .defaultValue(false) .withDescription( Enabling native aggregate function for hive dialect to use hash-agg strategy that can improve the aggregation performance.);在 HiveModule.java 中定义了原生聚合函数集合static final SetString BUILTIN_NATIVE_AGG_FUNC Collections.unmodifiableSet( new HashSet(Arrays.asList(sum, count, avg, min, max)));getFunctionDefinition中当开启该选项且函数名命中上述集合时会返回 Flink 原生实现见 HiveModule.java 与 getBuiltInNativeAggFunction对应实现类分别为HiveSumAggFunction→ HiveSumAggFunction.javaHiveCountAggFunction→ HiveCountAggFunction.javaHiveAverageAggFunction→ HiveAverageAggFunction.javaHiveMinAggFunction→ HiveMinAggFunction.javaHiveMaxAggFunction→ HiveMaxAggFunction.java这些类以 FlinkAggregateFunction的形式实现因而可以走 Flink 基于哈希的聚合算子执行路径替代 Hive 原实现只能使用 sort-based 聚合的限制。Hive 用户自定义函数UDFFlink 支持直接复用用户在 Hive 中已有的用户自定义函数UDF无需重写。支持的 UDF 类型UDFGenericUDFGenericUDTFUDAFGenericUDAFResolver2自动转换规则在查询规划与执行阶段Flink 会将 Hive 函数自动翻译为对应的 Flink 函数类型Hive 函数类型转换为 Flink 函数类型UDF/GenericUDFScalarFunction标量函数GenericUDTFTableFunction表值函数UDAF/GenericUDAFResolver2AggregateFunction聚合函数从源码结构看该转换由HiveFunctionDefinitionFactory负责HiveModule中所有 Hive 函数包括上述 UDF均通过它创建 Flink 的FunctionDefinition见 HiveModule.java。使用 Hive UDF 的前提条件要在 Flink 中使用 Hive 用户自定义函数需要满足两个条件设置 HiveCatalog将会话的当前 catalog 设置为由 Hive Metastore 支持的HiveCatalog且该 Metastore 中已注册对应的函数将包含该函数的 jar 加入 Flink classpath确保 Flink 运行时能够加载到函数实现类。关于 HiveCatalog 的完整配置方法可参考 Hive Catalog 文档。Hive UDF 使用实战假设我们在 Hive Metastore 中注册了如下三个 Hive 函数源码示例简单 UDF注册名为myudf/** * Test simple udf. Registered under name myudf */ public class TestHiveSimpleUDF extends UDF { public IntWritable evaluate(IntWritable i) { return new IntWritable(i.get()); } public Text evaluate(Text text) { return new Text(text.toString()); } }通用 UDF注册名为mygenericudf/** * Test generic udf. Registered under name mygenericudf */ public class TestHiveGenericUDF extends GenericUDF { Override public ObjectInspector initialize(ObjectInspector[] arguments) throws UDFArgumentException { checkArgument(arguments.length 2); checkArgument(arguments[1] instanceof ConstantObjectInspector); Object constant ((ConstantObjectInspector) arguments[1]).getWritableConstantValue(); checkArgument(constant instanceof IntWritable); checkArgument(((IntWritable) constant).get() 1); if (arguments[0] instanceof IntObjectInspector || arguments[0] instanceof StringObjectInspector) { return arguments[0]; } else { throw new RuntimeException(Not support argument: arguments[0]); } } Override public Object evaluate(DeferredObject[] arguments) throws HiveException { return arguments[0].get(); } Override public String getDisplayString(String[] children) { return TestHiveGenericUDF; } }UDTF注册名为myudtf/** * Test split udtf. Registered under name mygenericudtf */ public class TestHiveUDTF extends GenericUDTF { Override public StructObjectInspector initialize(ObjectInspector[] argOIs) throws UDFArgumentException { checkArgument(argOIs.length 2); // TEST for constant arguments checkArgument(argOIs[1] instanceof ConstantObjectInspector); Object constant ((ConstantObjectInspector) argOIs[1]).getWritableConstantValue(); checkArgument(constant instanceof IntWritable); checkArgument(((IntWritable) constant).get() 1); return ObjectInspectorFactory.getStandardStructObjectInspector( Collections.singletonList(col1), Collections.singletonList(PrimitiveObjectInspectorFactory.javaStringObjectInspector)); } Override public void process(Object[] args) throws HiveException { String str (String) args[0]; for (String s : str.split(,)) { forward(s); forward(s); } } Override public void close() { } }确认函数已在 Hive 中注册在 Hive CLI 中执行show functions可以看到这三个函数已经注册hive show functions; OK ...... mygenericudf myudf myudtf在 Flink SQL 中调用在 Flink SQL Client 中配合前面加载好的HiveModule或已设置为当前 catalog 的 HiveCatalog即可在 SQL 中同时使用这些 UDF包括与 lateral table 结合使用 UDTFFlink SQL select mygenericudf(myudf(name), 1) as a, mygenericudf(myudf(age), 1) as b, s from mysourcetable, lateral table(myudtf(name, 1)) as T(s);该示例同时展示了myudf(name)、myudf(age)以 UDF 方式处理字段再作为参数传给mygenericudfmygenericudf(..., 1)GenericUDF 的第二个参数是常量参数示例中校验了常量值为1lateral table(myudtf(name, 1)) as T(s)GenericUDTF 以横向表lateral table形式使用将name按逗号拆分并输出列s。总结Flink 通过三层次机制完整复用了 Hive 函数生态HiveModule把 Hive 内置函数注册为 Flink 系统函数借助模块解析顺序USE MODULES hive, core可控制同名函数的解析优先级源码层面通过黑名单机制保护 Flink 自身的窗口与解析函数原生 Hive 聚合函数table.exec.hive.native-agg-function.enabled默认false让sum/count/avg/min/max走 Flink 的 hash 聚合算子以提升性能但需注意其能力尚未与 Hive 内置实现完全对齐且目前需在加载HiveModule之前开启Hive UDF支持UDF/GenericUDF/GenericUDTF/UDAF/GenericUDAFResolver2在查询规划阶段自动转换为 Flink 的ScalarFunction/TableFunction/AggregateFunction前提是配置好承载函数的 HiveCatalog并将函数 jar 放入 Flink classpath。相关文档还可继续参考 Hive 连接器概述、Hive Catalog 与 Hive 读写以构建完整的 Hive 集成能力。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐Apache Spark SQL 集成 Hive UDF / UDAF / UDTF 完整指南Apache Spark SQL 集成 Hive UDF / UDAF / UDTF 完整指南 导读 本文以 Apache Spark 官方 SQL 参考文档大数据数据分析批处理流处理机器学习图计算Apache Arrow pyarrow.compute 计算函数库完整指南:从聚合、字符串函数到 UDF 注册Apache Arrow pyarrow.compute 计算函数库完整指南:从聚合、字符串函数到 UDF 注册 本篇技术文章以 Apache Arrow 官方大数据数据分析数据工程序列化深入拆解 webpack HarmonyES Module打包机制从 examples/harmony 看 import/export 编译与 import() 代码分割深入拆解 webpack HarmonyES Module打包机制从 examples/harmony 看 import/export 编译与 impor数据库OLAP数据仓库大数据湖仓一体数据分析上一篇ThinkLibrary人力资源招聘管理系统实战指南下一篇十条蛍LoRA合作伙伴生态系统建设与合作创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表