如果你想要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 函数它应用了自己的一组运行时优化。因此,在实践中,您必须凭经验测试性能,才能真正评估每种方法的好处。