【问题标题】:Apache Drill: Write general-purpose array_agg UDFApache Drill:编写通用的array_agg UDF
【发布时间】: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_agg UDF?
  • 如何使 UDF 与类型无关/通用?我真的必须为所有类型的每个 Nullable、Required 和 Repeated 版本实现一个完整的类吗?这是很多事情要做,而且相当乏味。没有办法处理与底层类型无关的 UDF 中的值吗? 我希望 Apache Drill 能够使用 Java 在这里提供的函数泛型类型、专门的函数重载和它们自己的类型系统的继承。我是否缺少有关如何执行此操作的内容?
  • 当我在聚合的 varchar 版本上使用 ORDER BY 时,如何解决 NULL 问题?
  • 如何解决我的地图/字典聚合为空列表的问题?
  • 是否有替代使用已弃用的ObjectHolder 的替代方法?

【问题讨论】:

    标签: java apache-drill


    【解决方案1】:

    为了回答您的问题,不幸的是,您遇到了 Drill Aggregate UDF API 的限制之一,即它只能返回简单的数据类型。1 Drill 解决这个问题将是一个很大的改进,但这是目前的状态。如果您有兴趣进一步讨论,请在 Drill 用户组和/或 slack 频道上启动一个线程。我不认为这是不可能的,但它需要对 Drill 内部进行一些修改。恕我直言,这非常值得,因为我想实现的其他一些 UDF 需要此功能。

    您问题的第二部分是如何使 UDF 类型不可知,再一次……您在 UDF API 中发现了另一个丑陋之处。 :-) 如果您在代码库中进行一些挖掘,您会发现大多数数学函数都有接受 FLOATINT 等的版本。

    关于 null 或空列表的聚合。我实际上在这里有一些好消息......目前这样做的方法是提供两个版本的函数,一个接受常规持有者,第二个接受可为空的持有者,如果输入为空,则返回一个空列表或映射。是的,这很糟糕,但另外一个好消息是我正在努力清理这个问题,希望很快会提交一份 PR,这样就不需要这样做了。

    关于 ObjectHolder,我编写了一个 median 函数,该函数使用一些堆栈来计算流式中位数,为此我使用了 ObjectHolder。我认为它会伴随我们一段时间,因为目前别无选择。

    我希望这能回答你的问题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-18
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多