
Apache Spark SQL 集成 Hive UDF / UDAF / UDTF 完整指南【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark导读本文以 Apache Spark 官方 SQL 参考文档 docs/sql-ref-functions-udf-hive.md 为骨架系统讲解 Spark SQL 如何复用 Hive 生态中三种用户自定义函数UDF标量函数、UDAF聚合函数与 UDTF表生成函数。通过本文你将掌握CREATE TEMPORARY FUNCTION注册机制、Hive 函数接口的选型逻辑UDFvsGenericUDF、UDAFvsGenericUDAFResolver以及 Spark SQL 在 Catalyst 执行引擎中的底层适配原理与性能注意事项。一、概述Spark SQL 对 Hive 函数的三种集成形态Spark SQL 原生支持集成 Hive UDF、UDAF 与 UDTF其行为语义与 Spark 自身的函数体系一一对应Hive UDFUser Defined Function单行输入、单行输出等价于 Spark SQL 标量函数Hive UDAFUser Defined Aggregate Function多行输入、返回单行聚合结果等价于 Spark 聚合函数Hive UDTFUser Defined Tabular Function单行输入、多行表状输出等价于 Spark 的Generator生成器典型如explode。使用流程分两步先在 Spark 中注册这些函数然后在 Spark SQL 查询中调用它们。注册动作本身不加载任何业务数据只是把Hive 类名 → Spark SQL 函数名的映射登记进当前会话或全局的 FunctionRegistry。底层适配HiveUDFExpressionBuilder从源码结构看Hive 函数被转换为 Catalyst 表达式的工作由 HiveSessionStateBuilder.scala 中的HiveUDFExpressionBuilder完成。它先尝试把注册的类当作 Spark 函数解析若抛出InvalidUDFClassException再按继承关系分派到对应的 Hive 包装表达式注册类的类型生成的 Catalyst 表达式继承org.apache.hadoop.hive.ql.exec.UDFHiveSimpleUDF实现org.apache.hadoop.hive.ql.udf.generic.GenericUDFHiveGenericUDF继承AbstractGenericUDAFResolverHiveUDAFFunctionisUDAFBridgeRequired false继承org.apache.hadoop.hive.ql.exec.UDAFHiveUDAFFunctionisUDAFBridgeRequired true经SparkGenericUDAFBridge桥接继承org.apache.hadoop.hive.ql.udf.generic.GenericUDTFHiveGenericUDTF这些表达式的实现集中在 hiveUDFs.scala底层求值器位于 hiveUDFEvaluators.scala。它们都实现了UserDefinedExpression接口因此可以与 Spark 原生 UDF 走同一套查询规划与代码生成管线。二、UDF注册与使用含性能提示Hive 提供两套 UDF 接口UDF简单 UDF通过反射调用具体的evaluate()方法要求参数类型与返回值类型匹配GenericUDF通用 UDF通过ObjectInspector描述输入输出支持更复杂的数据类型与类型推断Hive 内置函数大多基于此实现。示例 1使用 GenericUDFAbs基于 GenericUDFGenericUDFAbs是 Hive 内置求绝对值的 UDF。注册后即可像普通标量函数一样使用-- 注册 GenericUDFAbs 并在 Spark SQL 中使用。 -- 注意如果使用自己开发的函数需要把包含它的 JAR 加入 classpath -- 例如ADD JAR yourHiveUDF.jar; CREATE TEMPORARY FUNCTION testUDF AS org.apache.hadoop.hive.ql.udf.generic.GenericUDFAbs; SELECT * FROM t; ----- |value| ----- | -1.0| | 2.0| | -3.0| ----- SELECT testUDF(value) FROM t; -------------- |testUDF(value)| -------------- | 1.0| | 2.0| | 3.0| --------------示例 2使用 UDFSubstr基于 UDF 的简单 UDFUDFSubstr是 Hive 经典evaluate()风格 UDF-- 注册 UDFSubstr 并在 Spark SQL 中使用。 -- 注意如果返回类型与方法参数都使用 Java 原生类型primitive可以获得更好的性能。 -- 例如 UDFSubstr 的数据处理链路为 UTF8String - Text - String -- 使用原生类型可以省掉 UTF8String - Text 这一步转换。 CREATE TEMPORARY FUNCTION hive_substr AS org.apache.hadoop.hive.ql.udf.UDFSubstr; SELECT hive_substr(Spark SQL, 1, 5) AS value; ----- |value| ----- |Spark| -----性能关键点优先使用 Java 原生类型原文档特别强调当 Hive UDF 的返回值与方法参数使用 Java 原生类型primitive时可获得更优性能。原因在于 Spark SQL 内部字符串以UTF8String表示而 Hive 传统 UDF 默认以TextHadoop Writable交互。以UDFSubstr为例其处理链路为UTF8String - Text - String若参数/返回值使用原生StringSpark 可以省去UTF8String - Text的中间转换直接进入UTF8String - String路径。源码验证两类 UDF 的求值差异HiveSimpleUDFhiveUDFs.scala通过SparkDefaultUDFMethodResolver根据子表达式类型解析evaluate方法然后借助ConversionHelper做参数类型转换最后反射调用并 unwrap 返回值HiveGenericUDFhiveUDFs.scala调用GenericUDF.initialize(argumentInspectors)确定返回类型与ObjectInspector再通过DeferredObject惰性求值。值得注意 hiveUDFEvaluators.scala 中HiveGenericUDFEvaluator会对常参数 确定性的调用做常量折叠constant folding即若所有入参都是常量且 UDF 确定性编译期就直接求值一次避免每行重复调用。两类 UDF 均通过HiveUDFType注解deterministic与stateful属性判断是否确定性从而决定 Spark 能否做常量折叠、算子下推等优化见 hiveUDFEvaluators.scala。三、UDTF注册与使用Hive 的 UDTF 继承GenericUDTF典型实现是GenericUDTFExplode把数组/Map 拆成多行。在 Spark SQL 中UDTF 被包装为 Catalyst 的Generator表达式即HiveGenericUDTF见 hiveUDFs.scala。-- 注册 GenericUDTFExplode 并在 Spark SQL 中使用 CREATE TEMPORARY FUNCTION hiveUDTF AS org.apache.hadoop.hive.ql.udf.generic.GenericUDTFExplode; SELECT * FROM t; ------ | value| ------ |[1, 2]| |[3, 4]| ------ SELECT hiveUDTF(value) FROM t; --- |col| --- | 1| | 2| | 3| | 4| ---底层原理Collector 收集多行输出HiveGenericUDTF通过自定义UDTFCollector实现一行进、多行出它把输入行包装成 Hive 的 struct调用GenericUDTF.initialize(inputInspector)建立输出ObjectInspector再调用function.process(...)UDTF 每次collect()输出一行被UDTFCollector收集成InternalRow批量返回给上游算子见 hiveUDFs.scala。terminate()时调用function.close()以便 UDTF 收尾。一个重要的兼容性限制源码注释明确指出Spark 的Generator语义不允许在输入行之间维护状态因此依赖分区级状态或先close()再输出的 UDTF其行为与在 Hive 中并不完全一致。若函数需要在行间维持状态应当改写成用户自定义聚合UDAF因为聚合在分布式执行下有更清晰、更安全的状态语义见 hiveUDFs.scala。实践中绝大多数 UDTF如 explode、GenericUDTFParseUrlTuple不受此影响。四、UDAF注册与使用Hive 同样提供两套 UDAF 接口UDAF早期接口通过内部Evaluator的init/iterate/merge/terminate实现聚合GenericUDAFResolver现代接口按入参类型解析出GenericUDAFEvaluatorHive 内置聚合如 sum、avg多基于它。原文档示例使用继承GenericUDAFResolver的GenericUDAFSum-- 注册 GenericUDAFSum 并在 Spark SQL 中使用 CREATE TEMPORARY FUNCTION hiveUDAF AS org.apache.hadoop.hive.ql.udf.generic.GenericUDAFSum; SELECT * FROM t; -------- |key|value| -------- | a| 1| | a| 2| | b| 3| -------- SELECT key, hiveUDAF(value) FROM t GROUP BY key; ------------------ |key|hiveUDAF(value)| ------------------ | b| 3| | a| 3| ------------------注册后可像 Spark 原生聚合函数一样搭配GROUP BY、窗口函数使用。底层原理三种聚合状态格式HiveUDAFFunction继承 Spark 的TypedImperativeAggregate[HiveUDAFBuffer]见 hiveUDFs.scala其 ScalaDoc 详细说明了一个 Hive UDAF 在 Spark 中可能经历的三种聚合状态格式GenericUDAFEvaluator.AggregationBuffer实例Hive 原生聚合状态iterate()、merge()、terminatePartial()、terminate()均基于它可由ObjectInspector检查的 Java 对象Hive 用它产出可序列化的部分聚合结果以便跨节点 shuffleSpark SQL 值Spark 端序列化 Hive UDAF 状态时使用写入UnsafeRow后再取其字节数组。对应地Spark 内部通过wrap()/wrapperFor()Spark 值 → Hive 对象、unwrap()/unwrapperFor()Hive 对象 → Spark 值、terminatePartial()Buffer → 可检查对象完成状态互转AggregationBufferSerDehiveUDFs.scala则负责把部分聚合结果序列化/反序列化支撑 map 端局部聚合与 shuffle 后的全局归并。一个关键的兼容性 workaroundUPDATE 与 MERGE 混合某些 Hive UDAF 在其生命周期内不允许混合调用 UPDATE 与 MERGE例如 hash 聚合回退到 sort 聚合时Spark 可能对同一 UDAF 先 UPDATE 再 MERGE。为此Spark 在聚合缓冲区HiveUDAFBuffer(buf, canDoMerge)中跟踪是否已可 merge若检测到 UPDATE 后需要 MERGE会用 FINAL 模式的 evaluator 重建缓冲区再执行terminatePartial转换见 hiveUDFs.scala。这与测试套件 HiveUDFSuite.scala 中覆盖的多模式场景相互印证。五、函数注册与 JAR 加载要点注册语法临时函数CREATE TEMPORARY FUNCTION name AS com.example.MyUDF仅对当前会话生效持久函数可注册到 Hive Metastore供后续会话复用相关命令实现见 functions.scala 中的CreateFunctionCommand。原文档中两处注释需要重点记忆如果使用自己编程实现的函数必须先把包含它的 JAR 加入 classpath例如ADD JAR yourHiveUDF.jar;ADD JAR由HiveSessionResourceLoaderHiveSessionStateBuilder.scala处理它会把 JAR 同时注册到 Hive Client 与 Spark 会话资源加载器确保 Driver 与 Executor 的 classloader 都能找到类。参考测试 HiveUDFSuite.scala 中ADD JAR后再CREATE TEMPORARY FUNCTION udtf_count2 AS ...的完整流程。使用前置条件Spark 需以 Hive 支持方式构建sql/hive模块存在即常规-PhiveprofileUDF 类必须对 Spark 的 classloader 可见ADD JAR或--jars提交自定义函数建议同时考虑确定性声明HiveUDFType(deterministic true, stateful false)可让 Spark 对函数做常量折叠等优化。六、与 Spark 原生 UDF 的取舍集成 Hive 函数的核心价值是复用既有 Hive 资产存量 Hive 数仓中的 UDF/UDAF/UDTF 无需重写即可在 Spark SQL 中无缝调用迁移成本低。但也要注意性能Hive UDF 涉及 Hive ObjectInspector 体系与 Spark 内部行的双向转换通常慢于 Spark 原生Scala/Java 或 PythonUDF选型时应优先保证 Hive UDF 使用原生类型参数/返回值见上文性能提示语义差异UDTF 的行间状态语义、UDAF 的 UPDATE/MERGE 混合限制等与 Hive 存在细微差别分布式场景下应以 Spark 的 Catalyst 语义为准调试Hive UDF 抛出的异常会被 hiveUDFEvaluators.scala 包装为failedExecuteUserDefinedFunction错误信息含函数类名、入参与返回类型便于定位问题。七、延伸阅读官方 SQL 参考SQL 参考文档 与 UDF/UDAF 标量函数参考自定义标量函数sql-ref-functions-udf-scalar.mdHive 集成概览SQL 编程指南 中 Hive 表与函数相关章节核心实现源码hiveUDFs.scala、hiveUDFEvaluators.scala、HiveSessionStateBuilder.scala测试用例HiveUDFSuite.scala、HiveUDFDynamicLoadSuite.scala【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考