编写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。
- newBuffer()初始化缓冲区
- iterate():迭代处理输入数据
作用:处理每一条输入数据,更新缓冲区的中间状态。
参数说明:
- buffer:当前的缓冲区。
- args:输入参数数组,此处 args[0] 是待计算的整数。
逻辑:将输入的整数累加到 sum,同时将计数 count 加 1。
- merge():合并分布式计算的中间结果
作用:在分布式计算场景下,合并多个子任务的中间结果。
参数说明:
- buffer:主缓冲区,用于最终保存合并后的结果。
- partial:其他子任务的缓冲区,用于待合并的部分结果。
逻辑:将子任务的 sum 和 count 分别累加到主缓冲区中。
- terminate():计算最终结果
作用:在所有数据处理完成后,根据缓冲区的中间状态计算最终的平均值。
逻辑:用总累加和 sum 除以总计数 count,返回平均值;若没有数据(count=0),返回 0 避免除零异常。
步骤2:导出Jar包
编写调试完成代码后,通过IntelliJ IDEA工具编译代码并导出Jar包。
- 单击“Build”,参考下图分别单击“Build Project ”、“Build Artifacts...”分别对代码进行编译和打包。 图1 编译和打包
- 在弹出的对话框中单击Build。 图2 编译和打包

