Scala Spark从含嵌套序列的数据集创建DataFrame遇MatchError问题
解决Spark Scala中嵌套Seq转DataFrame的scala.MatchError问题
错误根源
报错出在map阶段的模式匹配逻辑:你遍历了每个元组的第二个元素(比如Seq(Seq(true, 5L), 0, 7.5)),并试图将单个元素匹配为(s: Seq[Any], i, d)三元组,但该序列里的第一个元素是Seq(true,5L)(单个嵌套序列),并非三元组,因此触发scala.MatchError。
修正方案
不需要遍历序列元素,直接解构嵌套序列的三层结构,将嵌套部分转为Row后,组合成与定义schema完全匹配的Row结构即可。
修正后的完整代码
import org.apache.spark.sql.types._ import org.apache.spark.sql.Row val data = Seq(("Java", Seq(Seq(true, 5L), 0, 7.5)), ("Python", Seq(Seq(true, 10L), 1, 8.5)), ("Scala", Seq(Seq(false, 8L), 2, 9.0))) // 精准解构嵌套序列,生成符合schema的Row val rdd = spark.sparkContext.parallelize(data).map { case (lang, Seq(userSeq: Seq[Any], difficulty: Int, avgReview: Double)) => val userRow = Row.fromSeq(userSeq) Row(lang, Row(userRow, difficulty, avgReview)) } val schema = StructType(Seq( StructField("language", StringType, true), StructField("stats", StructType(Seq( StructField("users", StructType(Seq( StructField("active", BooleanType, true), StructField("level", LongType, true) ))), StructField("difficulty", IntegerType, true), StructField("average_review", DoubleType, true) ))) )) val ds = spark.createDataFrame(rdd, schema) ds.show()
关键修正点
- 直接在模式匹配中解构
Seq(userSeq: Seq[Any], difficulty: Int, avgReview: Double),精准匹配每个元组第二个元素的三层结构 - 仅将嵌套的用户状态序列
userSeq转为Row,再组合成stats对应的Row,完全贴合schema定义的层级
运行输出
+--------+--------------------+ |language| stats| +--------+--------------------+ | Java| {[true,5], 0, 7.5} | | Python|{[true,10], 1, 8.5} | | Scala| {[false,8], 2, 9.0}| +--------+--------------------+
内容的提问来源于stack exchange,提问作者B-Brennan
相关产品推荐
相关产品推荐

