【问题标题】:How to convert nested object from rdd row to some custom object如何将嵌套对象从 rdd 行转换为一些自定义对象
【发布时间】:2020-01-05 19:28:20
【问题描述】:

我正在尝试学习一些 scala/spark 并尝试使用一些基本的 spark 集成示例进行练习。所以我的问题是我有一个本地运行的 Mongo 数据库。我正在提取一些数据并从中制作一个 rdd。 db 中的数据结构如下:

{
    "_id": 0,
    "name": "aimee Zank",
    "scores": [
        {
            "score": 1.463179736705023,
            "type": "exam"
        },
        {
            "score": 11.78273309957772,
            "type": "quiz"
        },
        {
            "score": 35.8740349954354,
            "type": "homework"
        }
    ]
}

这里有一些代码:

val conf: SparkConf = new SparkConf().setMaster("local[*]").setAppName("simple-app")
    val sparkSession = SparkSession.builder()
      .appName("example-spark-scala-read-and-write-from-mongo")
      .config(conf)
      .config("spark.mongodb.output.uri", "mongodb://sproot:12345@172.18.0.3:27017/spdb.students")
      .config("spark.mongodb.input.uri", "mongodb://sproot:12345@172.18.0.3:27017/spdb.students")
      .getOrCreate()

    // Reading Mongodb collection into a dataframe
    val df = MongoSpark.load(sparkSession)
    val dataRdd: RDD[Row] = df.rdd

    dataRdd.foreach(row => println(row.getValuesMap[Any](row.schema.fieldNames)))

上面的代码为我提供了这个:

Map(_id -> 0, name -> aimee Zank, scores -> WrappedArray([1.463179736705023,exam], [11.78273309957772,quiz], [35.8740349954354,homework]))
Map(_id -> 1, name -> Aurelia Menendez, scores -> WrappedArray([60.06045071030959,exam], [52.79790691903873,quiz], [71.76133439165544,homework]))

最后我把这些数据转换成:

case class Student(id: Long, name: String, scores: Scores)

case class Scores(@JsonProperty("scores") scores: List[Score])

case class Score (
                 @JsonProperty("score") score: Double,
                 @JsonProperty("type") scoreType: String
)

总结 - 问题是我无法将一些数据从 RDD 转换为 Student 对象。对我来说最有问题的地方是“得分”嵌套对象。 请帮助我了解应该如何做到这一点。

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    多玩了一下,最终得到了以下解决方案:

    object MainClass {
    
      def main(args: Array[String]): Unit = {
    
        val conf: SparkConf = new SparkConf().setMaster("local[*]").setAppName("simple-app")
        val sparkSession = SparkSession.builder()
          .appName("example-spark-scala-read-and-write-from-mongo")
          .config(conf)
          .config("spark.mongodb.output.uri", "mongodb://sproot:12345@172.18.0.3:27017/spdb.students")
          .config("spark.mongodb.input.uri", "mongodb://sproot:12345@172.18.0.3:27017/spdb.students")
          .getOrCreate()
    
        val objectMapper = new ObjectMapper()
        objectMapper.registerModule(DefaultScalaModule)
    
        // Reading Mongodb collection into a dataframe
        val df = MongoSpark.load(sparkSession)
        val dataRdd: RDD[Row] = df.rdd
    
        val students: List[Student] =
          dataRdd
            .collect()
            .map(row => Student(row.getInt(0), row.getString(1), createScoresObject(row))).toList
        println()
      }
    
      def createScoresObject(row: Row): Scores = {
        Scores(getAllScoresFromWrappedArray(row).map(x => Score(x.getDouble(0), x.getString(1))).toList)
      }
    
      def getAllScoresFromWrappedArray(row: Row): mutable.WrappedArray[GenericRowWithSchema] = {
        getScoresWrappedArray(row).map(x => x.asInstanceOf[GenericRowWithSchema])
      }
    
      def getScoresWrappedArray(row: Row): mutable.WrappedArray[AnyVal] = {
        row.getAs[mutable.WrappedArray[AnyVal]](2)
      }
    }
    
    case class Student(id: Long, name: String, scores: Scores)
    
    case class Scores(scores: List[Score])
    
    case class Score (score: Double, scoreType: String)
    

    但我很高兴知道是否有一些优雅的解决方案。

    【讨论】:

      猜你喜欢
      • 2022-12-18
      • 2021-05-28
      • 1970-01-01
      • 2021-11-21
      • 2011-01-15
      • 2016-04-03
      • 2021-07-27
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多