Scala Spark新手求教:将(String, List[(String, String)])转JSON对象
如何在Scala Spark中将你的数据转换为目标JSON结构
作为Scala Spark新手,别担心,我们一步步来搞定这个转换。首先得明确你的输入数据结构——看起来是一个包含(ID, List[(键, 值)])的RDD(或者Dataset),我们的目标是把每个ID对应的键值对列表转成嵌套对象,最终生成你要的JSON格式。
方法1:用RDD + Scala JSON库(适合小数据集)
如果你的数据量不大,可以先把数据收集到Driver端,用Scala的JSON库(比如Play JSON)来生成目标JSON:
步骤1:准备输入RDD
先模拟你的输入数据(如果是从数据源读取的,替换成对应的读取逻辑即可):
import org.apache.spark.SparkContext import org.apache.spark.SparkConf val conf = new SparkConf().setAppName("BookDataToJson").setMaster("local[*]") val sc = new SparkContext(conf) val inputRDD = sc.parallelize(Seq( ("ID001", List(("BookType", "[text]"), ("author", "xyz abc"), ("time", "01/12/2019[22:00] CST/PM"))), ("ID002", List(("BookType", "[text]"), ("author", "klj fgh"), ("time", "19/02/2019[12:00] CST/AM"))) ))
步骤2:将List转换为Map
每个ID对应的List[(String, String)]可以直接用toMap转成键值对Map:
val idToMapRDD = inputRDD.map { case (id, kvList) => (id, kvList.toMap) }
步骤3:用Play JSON生成最终JSON
首先需要在你的项目依赖里添加Play JSON(比如在build.sbt中):
libraryDependencies += "com.typesafe.play" %% "play-json" % "2.9.2"
然后序列化Map为JSON:
import play.api.libs.json._ // 把RDD的数据收集到Driver端的Scala Map val idToMap = idToMapRDD.collectAsMap() // 转成JSON并打印 val targetJson = Json.toJson(idToMap) println(targetJson.toString())
运行后就能得到你想要的JSON结构啦!
方法2:用DataFrame API(适合大数据集)
如果数据量很大,不适合收集到Driver端,就用Spark的DataFrame API来分布式处理:
步骤1:转换为DataFrame
import org.apache.spark.sql.SparkSession import spark.implicits._ val spark = SparkSession.builder().config(conf).getOrCreate() val df = inputRDD.toDF("id", "kv_list")
步骤2:将键值对列表转成Map列
用Spark内置的map_from_entries函数把数组形式的键值对转成Map:
import org.apache.spark.sql.functions._ val dfWithMap = df.withColumn("book_info", map_from_entries(col("kv_list")))
步骤3:合并所有ID的Map为一个大Map
我们需要把每个ID和对应的book_info合并成一个顶级Map,再转成JSON:
// 先把每行转成(id -> book_info)的单个Map val singleMapDF = dfWithMap.select(map(col("id"), col("book_info")).alias("single_map")) // 聚合所有单个Map为一个大Map val finalMapDF = singleMapDF.agg(collect_list("single_map").alias("map_list")) .withColumn("final_map", map_concat(col("map_list"): _*)) // 转成JSON字符串 val jsonResult = finalMapDF.select(to_json(col("final_map")).alias("target_json")) .collect()(0).getAs[String]("target_json") println(jsonResult)
这种方式全程在集群上处理,不会把大量数据拉到Driver端,适合生产环境的大数据场景。
需要注意的是,如果你的输入数据不是RDD而是其他格式(比如从文件读取的原始字符串),可能需要先做一步解析,把原始字符串转成(ID, List[(键, 值)])的结构,但核心的转换逻辑是一样的。
内容的提问来源于stack exchange,提问作者chris
相关产品推荐
相关产品推荐

