【发布时间】:2019-08-08 07:33:24
【问题描述】:
我有一个数据集,我正在提取并应用一个特定的模式,然后再写成 json。
我的测试数据集如下所示:
cityID|retailer|postcode
123|a1|1
123|s1|2
123|d1|3
124|a1|4
124|s1|5
124|d1|6
我想按城市 ID 分组。然后我应用以下模式并将其放入数据框中。然后我想把数据写成json。我的代码如下:
按城市 ID 分组
val rdd1 = cridf.rdd.map(x=>(x(0).toString, (x(1).toString, x(2).toString))).groupByKey()
将 RDD 映射到行
val final1 = rdd1.map(x=>Row(x._1,x._2.toList))
应用架构
val schema2 = new StructType()
.add("cityID", StringType)
.add("reads", ArrayType(new StructType()
.add("retailer", StringType)
.add("postcode", IntegerType)))
创建数据框
val parsedDF2 = spark.createDataFrame(final1, schema2)
写入 json 文件
parsedDF2.write.mode("overwrite")
.format("json")
.option("header", "false")
.save("/XXXX/json/testdata")
作业由于以下错误而中止:
java.lang.RuntimeException:编码时出错:
java.lang.RuntimeException: scala.Tuple2 不是结构模式的有效外部类型
【问题讨论】:
-
@JānisŠ。在 final1 x._2 中是零售商和邮政编码的列表
-
是的,我忽略了这一点。
标签: json scala dataframe apache-spark apache-spark-sql