Scala读取大CSV文件为RDD[Vector]时遇类型不匹配错误求助
解决Scala+Spark中CSV转RDD[Vector]的类型不匹配问题
兄弟,你碰到的这个类型不匹配问题,根源其实很简单:Scala里不带yield的for循环返回的是Unit类型,而你把这个塞到Seq里再转成RDD,自然就和你声明的RDD[Vector]类型对不上了。
问题分析
你原来的代码里,Seq( for(line <- Source.fromFile(file).getLines){ ... } ) 这段逻辑有核心问题:
- 不带
yield的for循环只是执行每行的处理逻辑,但不会返回任何有效结果,它的返回值是Unit(类似Java里的void) - 所以这个
Seq里其实只装了一个Unit对象,用sc.parallelize转完得到的是RDD[Unit],和你需要的RDD[Vector]类型完全不匹配,就触发了那个错误。
两种解决方案
方案1:修正原代码(适合小文件场景)
给for循环加上yield,让它把每次处理得到的Vector收集起来,生成一个Seq[Vector]:
import org.apache.spark.mllib.linalg.Vectors import org.apache.spark.rdd.RDD import scala.io.Source val file = "/data.csv" // 用yield让for循环返回Vector的迭代器,再转成Seq val vectorSeq = for(line <- Source.fromFile(file).getLines) yield { // 按逗号拆分每行,转成Double数组后生成稠密向量 val values = line.split(",").map(_.toDouble) Vectors.dense(values) } // 转成RDD val data: RDD[Vector] = sc.parallelize(vectorSeq.toSeq)
方案2:Spark原生分布式读取(适合大型CSV文件)
注意哦,你提到这是大型CSV文件,用Source.fromFile在Driver端读取整个文件其实不太合适——如果文件特别大,Driver的内存可能会扛不住。更推荐用Spark的textFileAPI分布式读取文件,效率更高:
import org.apache.spark.mllib.linalg.Vectors import org.apache.spark.rdd.RDD val file = "/data.csv" // 分布式读取文件,每个分区处理一部分行,直接映射成Vector val data: RDD[Vector] = sc.textFile(file).map { line => val values = line.split(",").map(_.toDouble) Vectors.dense(values) }
这个写法不需要在Driver端加载整个文件,Spark会自动把文件拆分到各个Executor节点处理,更适配大数据场景。
内容的提问来源于stack exchange,提问作者Tolga
相关产品推荐
相关产品推荐

