Spark DataFrame转JSON写入报Unable to find encoder for type Unit错如何解决
问题原因
- 核心报错原因:你在
map算子的函数体中只声明了details_1变量但没有返回该变量,整个匿名函数的返回值为Unit类型,Spark找不到对应Unit类型的编码器,因此抛出该错误。 - 嵌套Case Class设计冗余:你定义的多层嵌套Case Class会生成嵌套结构的JSON,和你需要的扁平JSON结构不符,完全不需要单独给每个字段包一层Case Class。
- 字段类型与取值错误:ID、Revenue字段都是数值类型,你统一用
r.getString()取值会出现类型转换错误;同时Scala中整数类型关键字是首字母大写的Int,你写的int本身也会编译报错。 - 字段名大小写不匹配:你定义的
details类中第二个字段是小写的country,最终生成的JSON键名也会是小写,和你需要的大写Country不符。
修正后可运行代码
import spark.implicits._ // 直接定义和目标JSON结构一致的扁平Case Class,字段名、类型完全匹配 case class RevenueInfo(ID: Int, Country: String, Revenue: Int) val DF = spark.sql("select ID,Country,Revenue from table") // map算子直接返回构造好的Case Class对象即可 val jsonDF = DF.map(r => { RevenueInfo( r.getInt(0), // ID是Int类型,用getInt取值 r.getString(1), // Country是字符串类型 r.getInt(2) // Revenue是Int类型 ) }) // 写出JSON jsonDF.repartition(1) .write .option("multiLine", "true") .json("C:/Desktop/output/revenue_output")
补充说明
- 上述代码写出的路径
C:/Desktop/output/revenue_output是文件夹,Spark会自动在该文件夹下生成以part-开头的JSON数据文件,不会直接生成名为revenue.json的单文件,如果你需要固定文件名,可以在写出完成后通过本地文件IO将生成的part文件重命名即可。 - 如果你数据量较大,不要使用
repartition(1),会将所有数据拉到同一个分区处理,容易出现内存溢出问题。 - 如果不需要JSON格式化换行展示,可以去掉
option("multiLine", "true")配置,默认生成的JSON是一行一个对象,更适合大数据场景下的解析。
内容的提问来源于stack exchange,提问作者mohan111
相关产品推荐
相关产品推荐

