【问题标题】:Spark Scala convert RDD with Case Class to simple RDDSpark Scala 将带有案例类的 RDD 转换为简单的 RDD
【发布时间】:2023-02-24 01:50:05
【问题描述】:
这可以:
case class trans(atm : String, num: Int)
val array = Array((20254552,"ATM",-5100), (20174649,"ATM",5120))
val rdd = sc.parallelize(array)
val rdd1 = rdd.map(x => (x._1, trans(x._2, x._3)))
如何再次转换回像 rdd 这样的简单 RDD?
例如。 rdd: org.apache.spark.rdd.RDD[(Int, String, Int)]
我可以做到这一点,当然:
val rdd2 = rdd1.mapValues(v => (v.atm, v.num)).map(x => (x._1, x._2._1, x._2._2))
但是如果班级有很大的记录怎么办?例如。动态地。
【问题讨论】:
标签:
scala
apache-spark
rdd
case-class
【解决方案1】:
不确定你想要的通用性,但在你的RDD[(Int, trans)]示例中,你可以使用trans伴随对象的unapply方法来将你的案例类展平为一个元组。
所以,如果你有你的设置:
case class trans(atm : String, num: Int)
val array = Array((20254552,"ATM",-5100), (20174649,"ATM",5120))
val rdd = sc.parallelize(array)
val rdd1 = rdd.map(x => (x._1, trans(x._2, x._3)))
您可以执行以下操作:
import shapeless.syntax.std.tuple._
val output = rdd1.map{
case (myInt, myTrans) => {
myInt +: trans.unapply(myTrans).get
}
}
output
res15: org.apache.spark.rdd.RDD[(Int, String, Int)]
我们正在导入 shapeless.syntax.std.tuple._ 以便能够从我们的 Int + 扁平元组(myInt +: trans.unapply(myTrans).get 操作)中创建一个元组。
【解决方案2】:
案例类方法“productIterator”可以帮助转换为数组:
case class trans(atm : String, num: Int)
val value = trans("ATM", 5120)
val rdd = spark.sparkContext.parallelize(Seq(value))
rdd
.map(_.productIterator.toArray)