Scala环境下使用Spark将多维数组转换为DataFrame报错如何解决?
报错原因
createDataFrame方法要求传入的集合元素必须是Product子类(比如元组、样例类实例),你定义的a是Array[Array[String]]类型,二维数组不满足该类型约束,因此触发类型不匹配报错。
现有代码的其他逻辑问题
在解决报错前你需要先修正原有代码的多个逻辑错误:
joinDf.orderBy("rectangle")是转换算子,属于懒执行,你没有将排序后的结果赋值给新变量,后续遍历的仍然是未排序的原始数据,排序逻辑完全不生效count和previous_r变量定义在for循环内部,每次循环都会被重置为初始值,你的计数逻辑根本无法正常累加- 调用
collect将全量数据拉到Driver节点本地统计,数据量稍大就会触发Driver内存溢出,完全违背分布式计算的设计原则
最优解决方案(推荐)
你的需求就是按矩形分组统计命中次数,直接用Spark原生的分组聚合API即可,不需要自己写循环逻辑,性能更高且不会有内存风险:
// 按rectangle分组统计次数,再按rectangle升序排序,直接得到目标DataFrame val resultDf = joinDf.groupBy("rectangle") .count() .orderBy("rectangle") .toDF("rectangle", "ST_contains_true_count")
该代码得到的resultDf就是你需要的、包含矩形名称和对应命中次数的DataFrame。
如果你坚持要修正原有数组转DF的写法
如果要让你的现有逻辑跑通,需要修改数组定义和元素类型:
- 将存储结果的数组改为存储元组(Tuple2属于Product子类,符合createDataFrame的参数要求)
- 把计数相关变量移到循环外部,避免每次循环重置
- 排序后的结果赋值给新变量再遍历
修改后的示例代码如下:
// 先拿到排序后的数据集 val sortedDf = joinDf.orderBy("rectangle") // 定义存储结果的数组,元素为(矩形名称, 命中次数)的元组 var resultArr = Array[(String, Long)]() var previous_r: Option[Double] = None var count: Long = 0 for (row <- sortedDf.rdd.collect) { val r = row.getAs[String]("rectangle").toDouble previous_r match { case None => previous_r = Some(r) count = 1 case Some(pre) if pre == r => count += 1 case Some(pre) => resultArr = resultArr :+ (pre.toString, count) previous_r = Some(r) count = 1 } } // 循环结束后把最后一个矩形的统计结果加进去 if (previous_r.isDefined) { resultArr = resultArr :+ (previous_r.get.toString, count) } // 转DataFrame val df = spark.createDataFrame(resultArr).toDF("rectangle", "ST_contains_true_count")
内容的提问来源于stack exchange,提问作者Tianna Wrona
相关产品推荐
相关产品推荐

