【问题标题】:Saving a RDD of key value pairs to a CSV file将键值对的 RDD 保存到 CSV 文件
【发布时间】:2019-12-22 07:47:18
【问题描述】:

我有一个键值对的 RDD,我想将其保存为 CSV 文件。

我编写了这段代码来从 HDFS 的一系列文件中获取 RDD。

val result = sc.sequenceFile[String,String](filenames)
val rdd_j= result.map(x => x._2)
rdd_j.take(1).foreach(println)

这给了我作为键值对的输出。下面是输出。

 {"lat":-37.676842,"lon":144.899414,"geoHash8":"r1r19m0s","adminRegionId":2344705 }

有很多这样的行。

现在我想将所有行保存到单个 CSV 中,其中键作为列,它们的值作为单元格值。此外,某些行中可能缺少某些键。请帮忙!

【问题讨论】:

  • 你试过什么?您要将其保存为本地文件还是 HDFS 上?您想如何处理缺失值?
  • 我尝试使用 saveAsTextFile 和 write.csv 将其保存为 csv。我想将缺少的键填充为 Csv 中的空值
  • 我想保存在hdfs上

标签: scala apache-spark rdd


【解决方案1】:

如果所有预期列都已知,则可以将数据转换为 DataFrame 并使用 'from_json' 函数提取:

val value = "{\"lat\":-37.676842,\"lon\":144.899414,\"geoHash8\":\"r1r19m0s\",\"adminRegionId\":2344705 }"
val rdd_j = sparkContext.parallelize(Seq(value))

// schema - other expected columns can be added here
val schema = StructType(
  Seq(
    StructField(name = "lat", dataType = DoubleType, nullable = true),
    StructField(name = "lon", dataType = DoubleType, nullable = true)
  )
)
// action
val df = rdd_j.toDF("value")
val result = df
  .withColumn("fromJson", from_json($"value", schema))
  .select($"fromJson.*")

result.show(false)

result.write.csv("outputPath")

输出:

+----------+----------+
|lat       |lon       |
+----------+----------+
|-37.676842|144.899414|
+----------+----------+

PS当架构未知时,可以使用简单的方法:

val result=spark.read.json(rdd_j.toDS())

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-11-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-24
    • 1970-01-01
    相关资源
    最近更新 更多