Spark Shell调用toDF报错:toDF不是org.apache.spark.rdd.RDD成员问题排查
问题分析与解决办法
错误原因
你的问题核心出在两个关键细节上:
- Case Class的作用域问题:你把
HService这个case class定义在了代码块({ ... })内部,在Spark Shell环境中,局部作用域内的case class无法被sqlContext.implicits._提供的隐式转换识别——而toDF方法正是依赖这些隐式转换,才能把RDD转换成DataFrame。 - 重复创建Spark上下文:Spark Shell启动时已经自动初始化了
sc(SparkContext)和sqlContext(Spark 1.x)/spark(SparkSession,Spark 2.x),你手动创建的上下文会和默认上下文冲突,导致隐式转换无法正常生效。
修正后的代码
针对Spark 1.x版本(你当前使用SQLContext的场景)
先在Shell的顶级作用域定义case class(不能放在代码块里):
case class HService( uhid:String, locationid:String, doctorid:String, billdate: String, servicename: String, servicequantity: String, starttime: String, endtime: String, servicetype: String, servicecategory: String, deptname: String )
然后直接编写业务逻辑,无需手动创建SparkConf、SparkContext和SQLContext:
// 直接使用Shell默认的sqlContext import sqlContext.implicits._ val hospitalDataText = sc.textFile("/home/training/Desktop/Data/services.csv") val header = hospitalDataText.first() val hospitalData = hospitalDataText.filter(_ != header) val hData = hospitalData.map(_.split(",")).map(p => HService( p(0),p(1),p(2),p(3),p(4), p(5),p(6),p(7),p(8),p(9),p(10) )) hData.take(4).foreach(println) val hosService = hData.toDF() hosService.registerTempTable("HService") val results = sqlContext.sql("SELECT doctorid, count(uhid) as visits FROM HService GROUP BY doctorid order by visits desc") results.collect().foreach(println)
针对Spark 2.x及以上版本(更推荐的写法)
Spark 2.x开始统一使用SparkSession,代码会更简洁:
case class HService( uhid:String, locationid:String, doctorid:String, billdate: String, servicename: String, servicequantity: String, starttime: String, endtime: String, servicetype: String, servicecategory: String, deptname: String ) // 使用Shell默认的SparkSession import spark.implicits._ val hospitalDataText = spark.read.textFile("/home/training/Desktop/Data/services.csv") val header = hospitalDataText.first() val hospitalData = hospitalDataText.filter(_ != header) val hData = hospitalData.map(_.split(",")).map(p => HService( p(0),p(1),p(2),p(3),p(4), p(5),p(6),p(7),p(8),p(9),p(10) )) val hosService = hData.toDF() hosService.createOrReplaceTempView("HService") val results = spark.sql("SELECT doctorid, count(uhid) as visits FROM HService GROUP BY doctorid order by visits desc") results.show() // 用show()比collect().foreach(println)更适配DataFrame的输出
额外注意事项
- 如果你的CSV存在字段内部包含逗号的情况(比如
servicename是"General, Checkup"),直接用split(",")会导致数组越界,建议直接用Spark的CSV数据源读取:// Spark 2.x+ 直接读取CSV并自动处理表头 val hosService = spark.read .option("header", "true") .option("quote", "\"") // 处理带引号的字段 .csv("/home/training/Desktop/Data/services.csv") .as[HService] - 在Spark Shell中,所有需要被隐式转换识别的类(比如case class)都必须放在顶级作用域,不能嵌套在方法或代码块里。
内容的提问来源于stack exchange,提问作者Bhaskar Das
相关产品推荐
相关产品推荐

