【发布时间】:2020-12-22 21:10:43
【问题描述】:
我想为 Apache Drill 创建一个 array_agg UDF,以便能够将一个组的所有值聚合到一个值列表中。
这应该适用于任何主要类型(必需、可选)和次要类型(varchar、dict、map、int 等)
但是,我觉得 Apache Drill 的 UDF API 并没有真正使用继承和泛型。每种类型都有自己的编写器和处理程序,它们不能被抽象为处理任何类型。例如,ValueHolder 接口似乎纯粹是装饰性的,不能用于将 UDF 与任何类型进行类型无关的挂钩。
我目前的实现
我尝试通过使用 Java 的反射来解决这个问题,这样我就可以使用 ListHolder 的写入函数,而与原始值的持有者无关。
但是,我随后遇到了@FunctionTemplate 注释的限制。
我无法为任何值创建通用 UDF 注释(我尝试使用接口 ValueHolder:@param ValueHolder input。
所以对我来说,支持不同类型的唯一方法似乎是为每种类型设置单独的类。但我什至不能抽象太多并处理任何@Param input,因为input 仅在其定义的类中可见(即特定于类型)。
我的实现基于https://issues.apache.org/jira/browse/DRILL-6963 并为必需和可选的varchars创建了以下两个类(首先如何统一?)
@FunctionTemplate(
name = "array_agg",
scope = FunctionScope.POINT_AGGREGATE,
nulls = NullHandling.INTERNAL
)
public static class VarChar_Agg implements DrillAggFunc {
@Param org.apache.drill.exec.expr.holders.VarCharHolder input;
@Workspace ObjectHolder agg;
@Output org.apache.drill.exec.vector.complex.writer.BaseWriter.ComplexWriter out;
@Override
public void setup() {
agg = new ObjectHolder();
}
@Override
public void reset() {
agg = new ObjectHolder();
}
@Override public void add() {
if (agg.obj == null) {
// Initialise list object for output
agg.obj = out.rootAsList();
}
org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter listWriter =
(org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter)agg.obj;
listWriter.varChar().write(input);
}
@Override
public void output() {
((org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter)agg.obj).endList();
}
}
@FunctionTemplate(
name = "array_agg",
scope = FunctionScope.POINT_AGGREGATE,
nulls = NullHandling.INTERNAL
)
public static class NullableVarChar_Agg implements DrillAggFunc {
@Param NullableVarCharHolder input;
@Workspace ObjectHolder agg;
@Output org.apache.drill.exec.vector.complex.writer.BaseWriter.ComplexWriter out;
@Override
public void setup() {
agg = new ObjectHolder();
}
@Override
public void reset() {
agg = new ObjectHolder();
}
@Override public void add() {
if (agg.obj == null) {
// Initialise list object for output
agg.obj = out.rootAsList();
}
if (input.isSet != 1) {
return;
}
org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter listWriter =
(org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter)agg.obj;
org.apache.drill.exec.expr.holders.VarCharHolder outHolder = new org.apache.drill.exec.expr.holders.VarCharHolder();
outHolder.start = input.start;
outHolder.end = input.end;
outHolder.buffer = input.buffer;
listWriter.varChar().write(outHolder);
}
@Override
public void output() {
((org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter)agg.obj).endList();
}
}
有趣的是,我无法导入 org.apache.drill.exec.vector.complex.writer.BaseWriter 以使整个事情变得更容易,因为这样 Apache Drill 将找不到它。
所以我必须把代码中org.apache.drill.exec.vector.complex.writer中所有东西的整个包路径都放在里面。
此外,我正在使用已废弃的 ObjectHolder。有更好的解决方案吗?
无论如何:这些工作到目前为止,例如用这个查询:
SELECT
MIN(tbl.`timestamp`) AS start_view,
MAX(tbl.`timestamp`) AS end_view,
array_agg(tbl.eventLabel) AS label_agg
FROM `dfs.root`.`/path/to/avro/folder` AS tbl
WHERE tbl.data.slug IS NOT NULL
GROUP BY tbl.data.slug
但是,当我使用 ORDER BY 时,我得到了这个:
org.apache.drill.common.exceptions.UserRemoteException: SYSTEM ERROR: UnsupportedOperationException: NULL
Fragment 0:0
此外,我尝试了更复杂的类型,即地图/字典。
有趣的是,当我打电话给SELECT sqlTypeOf(tbl.data) FROM tbl 时,我得到了 MAP。
但是当我编写 UDF 时,查询规划器抱怨没有 UDF array_agg 用于类型 dict。
不管怎样,我为dicts写了一个版本:
@FunctionTemplate(
name = "array_agg",
scope = FunctionScope.POINT_AGGREGATE,
nulls = NullHandling.INTERNAL
)
public static class Map_Agg implements DrillAggFunc {
@Param MapHolder input;
@Workspace ObjectHolder agg;
@Output org.apache.drill.exec.vector.complex.writer.BaseWriter.ComplexWriter out;
@Override
public void setup() {
agg = new ObjectHolder();
}
@Override
public void reset() {
agg = new ObjectHolder();
}
@Override public void add() {
if (agg.obj == null) {
// Initialise list object for output
agg.obj = out.rootAsList();
}
org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter listWriter =
(org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter) agg.obj;
//listWriter.copyReader(input.reader);
input.reader.copyAsValue(listWriter);
}
@Override
public void output() {
((org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter)agg.obj).endList();
}
}
@FunctionTemplate(
name = "array_agg",
scope = FunctionScope.POINT_AGGREGATE,
nulls = NullHandling.INTERNAL
)
public static class Dict_agg implements DrillAggFunc {
@Param DictHolder input;
@Workspace ObjectHolder agg;
@Output org.apache.drill.exec.vector.complex.writer.BaseWriter.ComplexWriter out;
@Override
public void setup() {
agg = new ObjectHolder();
}
@Override
public void reset() {
agg = new ObjectHolder();
}
@Override public void add() {
if (agg.obj == null) {
// Initialise list object for output
agg.obj = out.rootAsList();
}
org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter listWriter =
(org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter) agg.obj;
//listWriter.copyReader(input.reader);
input.reader.copyAsValue(listWriter);
}
@Override
public void output() {
((org.apache.drill.exec.vector.complex.writer.BaseWriter.ListWriter)agg.obj).endList();
}
}
但在这里,我在字段data_agg 中得到一个空列表用于我的查询:
SELECT
MIN(tbl.`timestamp`) AS start_view,
MAX(tbl.`timestamp`) AS end_view,
array_agg(tbl.data) AS data_agg
FROM `dfs.root`.`/path/to/avro/folder` AS tbl
GROUP BY tbl.data.viewSlag
问题总结
- 最重要的是:如何为 Apache Drill 创建
array_aggUDF? - 如何使 UDF 与类型无关/通用?我真的必须为所有类型的每个 Nullable、Required 和 Repeated 版本实现一个完整的类吗?这是很多事情要做,而且相当乏味。没有办法处理与底层类型无关的 UDF 中的值吗? 我希望 Apache Drill 能够使用 Java 在这里提供的函数泛型类型、专门的函数重载和它们自己的类型系统的继承。我是否缺少有关如何执行此操作的内容?
- 当我在聚合的 varchar 版本上使用
ORDER BY时,如何解决 NULL 问题? - 如何解决我的地图/字典聚合为空列表的问题?
- 是否有替代使用已弃用的
ObjectHolder的替代方法?
【问题讨论】:
标签: java apache-drill