# 在Spark SQL作业中使用UDAF
#### 操作场景
AI DataLake支持用户使用Hive UDAF（User Defined Aggregation Function，用户定义聚合函数）可对多行数据产生作用，通常与groupBy联合使用；等同于SQL中常用的SUM()，AVG()，也是聚合函数。
#### 约束限制
自定义函数中引用static类或接口时，必须要加上"try catch"异常捕获，否则可能会造成包冲突，导致函数功能异常。
#### 环境准备
在进行UDAF开发前，请准备以下开发环境。
表1UDAF开发环境 
| 准备项               | 说明                                                                                                                        |
|:---|:---|
| 操作系统              | Windows系统，支持Windows7以上版本。                                                                                                 |
| 安装JDK              | JDK使用21及以上版本（访问[Java官网](https://www.oracle.com/java/technologies/javase-downloads.html)）。                                   |
| 安装和配置IntelliJ IDEA | [IntelliJ IDEA](https://www.jetbrains.com/idea/)为进行应用开发的工具，版本要求使用2019.1或其2019.1往后的版本。                                     |
| 安装Maven             | 开发环境的基本配置（[下载](https://maven.apache.org/download.cgi)并[安装 Maven](https://maven.apache.org/install.html)）。用于项目管理，贯穿软件开发生命周期。 |
   
#### 开发流程
UDAF函数开发流程参考如下：
图1UDAF开发流程   
![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002671267250.png "点击放大")
表2开发流程说明 
| 序号  | 阶段                | 操作界面        | 说明                                                                                                                           |
|:---|:---|:---|:---|
| 1    | 新建Maven工程，配置pom文件 | IntelliJ IDEA | 参考[操作步骤]说明，编写UDAF函数代码。 |
| 2  | 编写UDAF函数代码        | IntelliJ IDEA | 参考[操作步骤]说明，编写UDAF函数代码。 |
| 3   | 调试，编译代码并导出Jar包   | IntelliJ IDEA | 参考[操作步骤]说明，编写UDAF函数代码。 |
| 4 | 上传Jar包到OBS         | OBS控制台       | 将生成的UDAF函数Jar包文件上传到OBS目录下。                                                                                                  |
| 5                                | 创建UDAF函数                                      | DataArts Studio控制台                        | 在DataArts Studio控制台的"脚本开发"界面，创建UDAF函数。                                                                                                                     |
| 6                                 | 验证和使用UDAF函数                                     | DataArts Studio控制台                          | 在DataArts Studio控制台的"脚本开发"界面，Spark SQL作业中使用创建的UDAF函数。                                                                                                       |
   
 #### 操作步骤
1. 新建Maven工程，配置pom文件。以下通过IntelliJ IDEA 2020.2工具操作演示。
   1. 打开IntelliJ IDEA，选择"File \> New \> Project"。
      图2新建Project   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002671267246.png "点击放大") 
   
   2. 选择Maven，Project SDK选择1.8，单击"Next"。
      图3配置Project SDK   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002701186819.png "点击放大") 
   
   3. 定义样例工程名和配置样例工程存储路径，单击"Create"，下一步单击弹窗中的"Finish"完成工程创建。
      图4完成Project创建   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002671267248.png "点击放大") 
   
   4. 在pom.xml文件中添加如下配置。
      ```
      <dependencies> 
               <dependency> 
                   <groupId>org.apache.hive</groupId> 
                   <artifactId>hive-exec</artifactId> 
                   <version>1.2.1</version> 
               </dependency> 
       </dependencies>
      ```
      图5pom文件中添加配置   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002671267252.png "点击放大") 
   
   5. 在工程路径的"src \> main \> java"文件夹上鼠标右键，选择"New \> Package"，新建Package和类文件。
      Package根据需要定义，本示例定义为："com.aidatalake.demo"
      图6新建Package   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002671427096.png "点击放大")
      在包路径下新建Java Class文件，本示例定义为：AvgFilterUDAFDemo。
      图7创建类   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002701306895.png "点击放大") 
    
2. 编写UDAF函数代码。UDAF函数实现，主要注意以下几点：
   - 自定义UDAF需要继承org.apache.hadoop.hive.ql.exec.UDAF和org.apache.hadoop.hive.ql.exec.UDAFEvaluator类。函数类需要继承UDAF类，计算类Evaluator实现UDAFEvaluator接口。
   
   - Evaluator需要实现UDAFEvaluator的**init** 、**iterate** 、**terminatePartial** 、**merge** 、**terminate** 这几个函数。
     - **init**函数实现接口UDAFEvaluator的init函数。
     
     - **iterate**接收传入的参数，并进行内部的迭代。
     
     - **terminatePartial**无参数，其为iterate函数遍历结束后，返回遍历得到的数据，terminatePartial类似于 hadoop的Combiner。
     
     - **merge**接收terminatePartial的返回结果。
     
     - **terminate**返回最终的聚集函数结果。
     
     
     详细UDAF函数实现，可以参考如下样例代码：
     ```
     package com.aidatalake.demo;
      
     import org.apache.hadoop.hive.ql.exec.UDAF;
     import org.apache.hadoop.hive.ql.exec.UDAFEvaluator;
      
     /***
      * @jdk jdk1.8.0
      * @version 1.0
      ***/
     public class AvgFilterUDAFDemo extends UDAF {
      
         /**
          * 定义静态内部类AvgFilter
          */
         public static class PartialResult
         {
             public Long sum;
         }
      
         public static class VarianceEvaluator implements UDAFEvaluator {
      
             //初始化PartialResult对象
             private AvgFilterUDAFDemo.PartialResult partial;
      
             //创建VarianceEvaluator无参构造函数
             public VarianceEvaluator(){
      
                 this.partial = new AvgFilterUDAFDemo.PartialResult();
      
                 init();
             }
      
             /**
              * init函数类似于构造函数，用于UDAF的初始化
              */
             @Override
             public void init() {
      
                 //设置sum初始值
                 this.partial.sum = 0L;
             }
      
             /**
              * iterate接收传入的参数，并进行内部的轮转。
              * @param x
              * @return
              */
             public void iterate(Long x) {
                 if (x == null) {
                     return;
                 }
                 AvgFilterUDAFDemo.PartialResult tmp9_6 = this.partial;
                 tmp9_6.sum = tmp9_6.sum | x;
             }
      
             /**
              * terminatePartial无参数，其为iterate函数遍历结束后，返回轮转数据，
              * terminatePartial类似于hadoop的Combiner
              * @return
              */
             public AvgFilterUDAFDemo.PartialResult terminatePartial()
             {
                 return this.partial;
             }
      
             /**
              * merge接收terminatePartial的返回结果，进行数据merge操作
              * @param
              * @return
              */
             public void merge(AvgFilterUDAFDemo.PartialResult pr)
             {
                 if (pr == null) {
                     return;
                 }
                 AvgFilterUDAFDemo.PartialResult tmp9_6 = this.partial;
                 tmp9_6.sum = tmp9_6.sum | pr.sum;
             }
      
             /**
              * terminate返回最终的聚集函数结果
              * @return
              */
             public Long terminate()
             {
                 if (this.partial.sum == null) {
                     return 0L;
                 }
                 return this.partial.sum;
             }
         }
     }
     ```
     图8编写UDAF函数代码   
     ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002701306899.png "点击放大") 
    
3. 编写调试完成代码后，通过IntelliJ IDEA工具编译代码并导出Jar包。
   1. 单击工具右侧的"Maven"，参考下图分别单击"clean"、"compile"对代码进行编译。 编译成功后，单击"package"对代码进行打包。
      图9导出jar包   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002701306897.png "点击放大") 
   
   2. 打包成功后，生成的Jar包会放到target目录下，以备后用。本示例将会生成到："D:\\AIDataLakeTest\\MyUDAF\\target"下名为"MyUDAF-1.0-SNAPSHOT.jar"。
    
4. 登录OBS控制台，将生成的Jar包文件上传到OBS路径下。
   ![](https://support.huaweicloud.com/devg-aidatalake/public_sys-resources/note_3.0-zh-cn.png)
   Jar包文件上传的OBS桶所在的区域需与AI DataLake的工作空间区域相同，不可跨区域执行操作。
   
5. 创建UDAF函数。
   1. 在[AI DataLake管理控制台](https://console.huaweicloud.com/fabric/)左下角，单击"DataArts Studio"。
   
   2. 在"脚本开发"页面，新建Serverless Spark SQL脚本后，选择端点，Catalog和数据库。
      图10新建Serverless Spark SQL脚本   
      ![](https://support.huaweicloud.com/devg-aidatalake/zh-cn_image_0000002671267254.png) 
   
   3. 在SQL编辑区域，输入创建UDAF函数的命令。
      ```
      CREATE FUNCTION AvgFilterUDAFDemo AS 'com.aidatalake.demo.AvgFilterUDAFDemo' using jar 'obs://aidatalake-test-obs01/MyUDAF-1.0-SNAPSHOT.jar';
      ```
      或
      ```
      CREATE OR REPLACE FUNCTION AvgFilterUDAFDemo AS 'com.aidatalake.demo.AvgFilterUDAFDemo' using jar 'obs://aidatalake-test-obs01/MyUDAF-1.0-SNAPSHOT.jar';
      ```
      ![](https://support.huaweicloud.com/devg-aidatalake/public_sys-resources/note_3.0-zh-cn.png)
      如果该用户开启了自定义函数热加载功能，注册语句会发生变化。
      
   
   4. 单击"运行"。
    
6. 使用UDAF函数。 在查询语句中使用创建的UDAF函数:
   ```
   select AvgFilterUDAFDemo(real_stock_rate) AS show_rate FROM dw_ad_estimate_real_stock_rate limit 1000;
   ```
   
7. （可选）删除UDAF函数。 如果不再使用UDAF函数，可执行以下语句删除该函数:
   ```
   Drop FUNCTION AvgFilterUDAFDemo;
   ```
   
 
