【问题标题】:Is it possible to unit test Spark UDAFs?是否可以对 Spark UDAF 进行单元测试?
【发布时间】:2017-09-27 22:19:24
【问题描述】:

Spark UDAF 要求您实现多种方法,特别是 def update(buffer: MutableAggregationBuffer, input: Row): Unitdef merge(buffer1: MutableAggregationBuffer, buffer2: Row): Unit

假设我的测试中有一个 UDAF X、4 行 (r0, r1, r2, r3) 和两个聚合缓冲区 A, B。 我想看看这段代码是否产生了预期的结果:

X.update(A, r0)
X.update(A, r1)
X.update(B, r2)
X.update(B, r3)
X.merge(A, B)
X.evaluate(A)

与仅使用一个缓冲区对 4 行中的每一行调用 X.update 相同:

X.update(A, r0)
X.update(A, r1)
X.update(A, r2)
X.update(A, r3)
X.evaluate(A)

这样可以测试两种方法的正确性。 但是,我不知道如何编写这样的测试:用户代码似乎无法实例化MutableAggregationBuffer 的任何实现。

如果我只是从我的 4 行中创建一个 DF,并尝试使用groupBy().agg(...) 来调用我的 UDAF,Spark 甚至不会尝试以这种特定方式合并它们 - 因为它的行数很少,它不需要。

【问题讨论】:

    标签: scala unit-testing apache-spark apache-spark-sql user-defined-functions


    【解决方案1】:

    MutableAggregationBuffer 只是一个抽象类。您可以轻松创建自己的实现,例如:

    import org.apache.spark.sql.expressions._
    
    class DummyBuffer(init: Array[Any]) extends MutableAggregationBuffer {
      val values: Array[Any] = init
      def update(i: Int, value: Any) = values(i) = value
      def get(i: Int): Any = values(i)
      def length: Int = init.size
      def copy() = new DummyBuffer(values)
    }
    

    它不会取代“真实的东西”,但对于简单的测试场景来说应该足够了。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-03-16
      • 2010-09-08
      • 2022-11-18
      • 1970-01-01
      • 1970-01-01
      • 2011-11-15
      • 2020-02-21
      • 1970-01-01
      相关资源
      最近更新 更多