【问题标题】: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)
      

      【讨论】:

      • 我会尝试,但我看不到您身边的案例类使用情况。
      • 我的例子中第一行使用了你的案例类“trans”。
      • 是的,但它不是案例类
      • 答案已更新,案例类已添加
      • 好的,我今晚会尝试,虽然不是那么通用。让我看看。
      猜你喜欢
      • 1970-01-01
      • 2021-09-27
      • 2017-06-13
      • 2014-11-28
      • 2018-03-05
      • 1970-01-01
      • 2017-05-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多