【问题标题】:Spark Dataset select with typedcolumn使用 typedcolumn 选择 Spark 数据集
【发布时间】:2016-12-03 04:31:56
【问题描述】:

查看 spark DataSet 上的select() 函数,有各种生成的函数签名:

(c1: TypedColumn[MyClass, U1],c2: TypedColumn[MyClass, U2] ....)

这似乎暗示我应该能够直接引用 MyClass 的成员并且是类型安全的,但我不确定如何......

ds.select("member") 当然有效.. 似乎ds.select(_.member) 也可能以某种方式有效?

【问题讨论】:

    标签: scala apache-spark apache-spark-dataset


    【解决方案1】:

    如果你想要ds.select(_.member) 的等价物,只需使用map

    case class MyClass(member: MyMember, foo: A, bar: B)
    val ds: DataSet[MyClass] = ???
    val members: DataSet[MyMember] = ds.map(_.member)
    

    编辑:不使用map的论据。

    一种更高效的方法是通过投影,根本不使用map。您失去了编译时类型检查,但作为交换,Catalyst 查询引擎有机会做一些更优化的事情。正如@Sim 在下面的评论中所暗示的那样,主要的优化中心是不需要将MyClass 的全部内容从钨内存空间反序列化到JVM 堆内存中——只是为了调用访问器——然后序列化@ 的结果987654327@回到钨。

    为了举一个更具体的例子,让我们重新定义我们的数据模型,如下所示:

      // Make sure these are not nested classes 
      // (i.e. in a top level compilation units).
      case class MyMember(something: Double)
      case class MyClass(member: MyMember, foo: Int, bar: String)
    

    这些必须是case 类,以便SQLImplicits.newProductEncoder[T <: Product] 可以为我们提供Dataset[T] API 所需的隐式Encoder[MyClass]

    现在我们可以让上面的例子更具体:

      val ds: Dataset[MyClass] = Seq(MyClass(MyMember(1.0), 2, "three")).toDS()
      val membersMapped: Dataset[Double] = ds.map(_.member.something)
    

    要查看幕后发生的事情,我们使用explain() 方法:

    membersMapped.explain()
    
    == Physical Plan ==
    *(1) SerializeFromObject [input[0, double, false] AS value#19]
    +- *(1) MapElements <function1>, obj#18: double
       +- *(1) DeserializeToObject newInstance(class MyClass), obj#17: MyClass
          +- LocalTableScan [member#12, foo#13, bar#14]
    

    这使得与 Tungsten 之间的序列化非常明显。

    让我们使用投影[^1]得到相同的值:

    val ds2: Dataset[Double] = ds.select($"member.something".as[Double])
    ds2.explain()
    
    == Physical Plan ==
    LocalTableScan [something#25]
    

    就是这样!一步[^2]。除了将MyClass编码到原始Dataset之外,没有序列化。

    [^1]:投影定义为$"member.something"而不是$"value.member.something"的原因与Catalyst自动投影单个列DataFrame的成员有关。

    [^2]:公平地说,第一个物理计划中的步骤旁边的* 表示它们将由WholeStageCodegenExec 实现,因此这些步骤成为单个即时编译的 JVM 函数它应用了自己的一组运行时优化。因此,在实践中,您必须凭经验测试性能,才能真正评估每种方法的好处。

    【讨论】:

    • 请注意,由于 Spark 内部使用原始二进制数据而不是 Scala 类型,因此会有数据转换成本。
    • 在这种情况下使用 Dataset 有什么好处?它只是在性能上权衡类型安全吗?我完全不明白 Dataset 什么时候会有用!
    • 大多数时候你只想使用一个数据框。有时,为了与其他函数的互操作性,您可能希望进入 DataSet 空间以便能够在不创建 UDF 的情况下调用 mapflatMap 等。或其他一些侧面案例。
    【解决方案2】:

    select 的Scala DSL 中,有很多方法可以识别Column

    • 来自符号:'name
    • 来自字符串:$"name"col(name)
    • 来自表达式:expr("nvl(name, 'unknown') as renamed")

    要从Column 获取TypedColumn,您只需使用myCol.as[T]

    例如:ds.select(col("name").as[String])

    【讨论】:

    • 这个答案是正确的,但是请注意 this as[T] 不是类型安全的,因此如果假设错误类型,它可能会在 RT 中爆炸。
    • 好点。为了获得编译器的最大帮助,您必须完全切换到 Scala 类型,例如,ds.as[T].map { t: T =&gt; ... }。请注意,由于 Spark 在内部使用原始二进制数据而不是 Scala 类型,因此会有数据转换成本。
    猜你喜欢
    • 2017-10-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-11-09
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多