更新时间:2026-07-28 GMT+08:00
分享

编写UDAF函数代码并导出Jar包

步骤1:编写UDAF函数代码

UDAF函数实现,主要注意以下几点:

  • 自定义UDAF需继承UDAFHandler类,在类内编写UDAF函数。并实现如下方法:
    public abstract Writable newBuffer();
    public abstract void iterate(Writable buffer, Object[] args) throws FuncException;
    public abstract void merge(Writable buffer, Writable partial) throws FuncException;
    public abstract Object terminate(Writable buffer) throws FuncException;
  • 在方法上添加@TypeResolve注解,用于表明函数的入参和出参。
    入参和出参通过'->'分隔。入参和出参类型(resolveType)以及对应的SQL和Java数据类型请参考表1
    表1 resolveType类型

    SQL type

    Java type

    resolveType

    BooleanType

    Boolean.class

    BOOLEAN

    BooleanType

    Boolean.class

    BOOLEAN

    ByteType

    Byte.class

    BYTE

    ShortType

    Short.class

    SHORT

    IntegerType

    Integer.class

    INT

    LongType

    Long.class

    LONG

    FloatType

    Float.class

    FLOAT

    DoubleType

    Double.class

    DOUBLE

    StringType

    String.class

    STRING

    BinaryType

    byte[].class

    BINARY

    DecimalType.Fixed(precision, scale)

    BigDecimal.class

    Decimal

    DateType

    LocalDate.class

    DATE

    TimestampType

    ZonedDateTime.class

    TIMESTAMP

    TimestampNTZType

    LocalDateTime.class

    TIMESTAMP_NTZ

    YearMonthIntervalType

    Period.class

    INTERVAL_YEAR_MONTH

    DayTimeIntervalType

    Duration.class

    INTERVAL_DAY_TIME

    List

    List.class

    LIST

    Map

    Map.class

    MAP

    Struct

    Map.class

    STRUCT

  • 添加@deterministic注解:

    表明函数是确定性的,即在给定相同的输入时,总是返回相同的结果。

    详细UDAF函数实现,可以参考如下样例代码:

    • 功能描述:

      本函数用于计算数据的平均值,设计流程为:“初始化 > 迭代计算 > 合并结果 > 最终计算”。

    • 示例代码
      package com.huawei.demo;
      
      import com.huawei.spark.function.FuncException;
      import com.huawei.spark.function.handlers.UDAFHandler;
      import com.huawei.spark.function.requests.Writable;
      import com.huawei.spark.function.types.TypeResolve;
      
      import java.io.IOException;
      import java.io.ObjectInput;
      import java.io.ObjectOutput;
      
      // get average
      @TypeResolve("int->double")
      public class UDAFDemo extends UDAFHandler {
          private static class Buffer implements Writable {
              private double sum = 0.d;
      
              private int count = 0;
      
              public Buffer() {
              }
      
              @Override
              public void writeExternal(ObjectOutput out) throws IOException {
                  out.writeDouble(sum);
                  out.writeInt(count);
              }
      
              @Override
              public void readExternal(ObjectInput in) throws IOException {
                  sum = in.readDouble();
                  count = in.readInt();
              }
          }
      
          public Writable newBuffer() {
              return new com.huawei.demo.UDAFDemo.Buffer();
          }
      
          public void iterate(Writable buffer, Object[] args) throws FuncException {
              com.huawei.demo.UDAFDemo.Buffer buf = (com.huawei.demo.UDAFDemo.Buffer) buffer;
              Integer arg = (Integer) args[0];
              buf.sum += arg;
              buf.count++;
          }
      
          public void merge(Writable buffer, Writable partial) throws FuncException {
              com.huawei.demo.UDAFDemo.Buffer buf = (com.huawei.demo.UDAFDemo.Buffer) buffer;
              com.huawei.demo.UDAFDemo.Buffer par = (com.huawei.demo.UDAFDemo.Buffer) partial;
              buf.sum += par.sum;
              buf.count += par.count;
          }
      
          public Object terminate(Writable buffer) throws FuncException {
              com.huawei.demo.UDAFDemo.Buffer buf = (com.huawei.demo.UDAFDemo.Buffer) buffer;
              if (buf.count == 0) {
                  return 0;
              }
              return buf.sum / buf.count;
          }
      }

    package com.huawei.demo;
    
    import com.huawei.spark.function.FuncException;
    import com.huawei.spark.function.handlers.UDAFHandler;
    import com.huawei.spark.function.requests.Writable;
    import com.huawei.spark.function.types.TypeResolve;
    
    import java.io.IOException;
    import java.io.ObjectInput;
    import java.io.ObjectOutput;
    
    @TypeResolve("long->double")
    public class UDAFDemo extends UDAFHandler {
        private static class Buffer implements Writable {
    
            private double sum = 0;
    
            private int count = 0;
    
            public Buffer() {
            }
    
            @Override
            public void writeExternal(ObjectOutput out) throws IOException {
                out.writeDouble(sum);
                out.writeInt(count);
            }
    
            @Override
            public void readExternal(ObjectInput in) throws IOException {
                sum = in.readDouble();
                count = in.readInt();
            }
        }
    
        public Writable newBuffer() {
            return new Buffer();
        }
    
        public void iterate(Writable buffer, Object[] args) throws FuncException {
            Buffer buf = (Buffer) buffer;
            if(null != args[0]) {
                Long arg = (Long) args[0];
                buf.sum += arg;
                buf.count++;
            }
        }
    
        public void merge(Writable buffer, Writable partial) throws FuncException {
            Buffer buf = (Buffer) buffer;
            Buffer par = (Buffer) partial;
            buf.sum += par.sum;
            buf.count += par.count;
        }
    
        public Object terminate(Writable buffer) throws FuncException {
            Buffer buf = (Buffer) buffer;
            if (buf.count == 0) {
                return 0;
            }
            return buf.sum / buf.count;
        }
    }
  • 详细说明:
    • newBuffer()初始化缓冲区

      Buffer 类实现了 Writable 接口,用于存储聚合过程中的中间结果,主要作用是暂存计算过程中的中间值,支持序列化和反序列化。

      包含两个核心变量:

      • sum:累加总和。
      • count:累加计数,记录参与计算的元素个数。

      实现了两个序列化方法:

      • writeExternal:将 sum 和 count 写入输出流。
      • readExternal:从输入流读取 sum 和 count。
  • iterate():迭代处理输入数据

    作用:处理每一条输入数据,更新缓冲区的中间状态。

    参数说明:

    • buffer:当前的缓冲区。
    • args:输入参数数组,此处 args[0] 是待计算的整数。

    逻辑:将输入的整数累加到 sum,同时将计数 count 加 1。

  • merge():合并分布式计算的中间结果

    作用:在分布式计算场景下,合并多个子任务的中间结果。

    参数说明:

    • buffer:主缓冲区,用于最终保存合并后的结果。
    • partial:其他子任务的缓冲区,用于待合并的部分结果。

    逻辑:将子任务的 sum 和 count 分别累加到主缓冲区中。

  • terminate():计算最终结果

    作用:在所有数据处理完成后,根据缓冲区的中间状态计算最终的平均值。

    逻辑:用总累加和 sum 除以总计数 count,返回平均值;若没有数据(count=0),返回 0 避免除零异常。

步骤2:导出Jar包

编写调试完成代码后,通过IntelliJ IDEA工具编译代码并导出Jar包。

  1. 单击“Build”,参考下图分别单击“Build Project ”、“Build Artifacts...”分别对代码进行编译和打包。
    图1 编译和打包
  1. 在弹出的对话框中单击Build。
    图2 编译和打包
  1. 完成上述步骤之后,可以在项目目录当中看到out文件夹,展开out目录。

    其中的MyRemoteUDAF.jar就是之后要上传到函数工作流FunctionGraph当中的jar包。

    图3 查看MyRemoteUDAF.jar
  2. 右键单击MyRemoteUDAF.jar,选择Copy Path/Reference,直接复制文件所在目录,方便之后上传JAR包使用。

    本示例将会生成到:“D:\UDAFTest\MyRemoteUDAF\out\artifacts\MyRemoteUDAF_jar\”目录下,文件名为“MyRemoteUDAF.jar”。

    图4 复制文件所在目录

相关文档