【发布时间】:2018-06-28 13:46:54
【问题描述】:
您好,我有一个 RDD,它基本上是在读取 CSV 文件后制作的。 我已经定义了一个方法,它基本上根据输入参数将 rdd 的行映射到不同的案例类。
返回的RDD需要转成dataframe 当我尝试运行相同时,我得到以下错误。
定义的方法是
case class Australiafile1(sectionName: String, profitCentre: String, valueAgainst: String, Status: String)
case class Australiafile2(sectionName: String, profitCentre: String)
case class defaultclass(error: String)
def mapper(line: String, recordLayoutClassToBeUsed: String) = {
val fields = line.split(",")
var outclass = recordLayoutClassToBeUsed match {
case ("Australiafile1") => Australiafile1(fields(0), fields(1), fields(2), fields(3))
case ("Australiafile2") => Australiafile2(fields(0), fields(1))
}
outclass
}
该方法的输出用于创建如下数据框
val inputlines = spark.sparkContext.textFile(inputFile).cache().mapPartitionsWithIndex { (idx, lines) => if (idx == 0) lines.drop(numberOfLinesToBeRemoved.toInt) else lines }.cache()
val records = inputlines.filter(x => !x.isEmpty).filter(x => x.split(",").length > 0).map(lines => mapper(lines, recordLayoutClassToBeUsed))
import spark.implicits._
val recordsDS = records.toDF()
recordsDS.createTempView("recordtable")
val output = spark.sql("select * from recordtable").toDF()
output.write.option("delimiter", "|").option("header", "false").mode("overwrite").csv(outputFile)
收到的错误如下
线程“main”中的异常 java.lang.NoClassDefFoundError: 没有找到对应于具有 Serializable 的 Product 的 Java 类 在 scala.reflect.runtime.JavaMirrors$JavaMirror.typeToJavaClass(JavaMirrors.scala:1300) 在 scala.reflect.runtime.JavaMirrors$JavaMirror.runtimeClass(JavaMirrors.scala:192) 在 scala.reflect.runtime.JavaMirrors$JavaMirror.runtimeClass(JavaMirrors.scala:54) 在 org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$.apply(ExpressionEncoder.scala:60) 在 org.apache.spark.sql.Encoders$.product(Encoders.scala:275) 在 org.apache.spark.sql.LowPrioritySQLImplicits$class.newProductEncoder(SQLImplicits.scala:233) 在 org.apache.spark.sql.SQLImplicits.newProductEncoder(SQLImplicits.scala:33)
你能告诉我这有什么问题吗,我该如何克服这个问题?
【问题讨论】:
-
案例类应该写在主类之外
-
它们仅在 main 方法之外。此外,如果我只从映射器方法返回一个案例类对象而不匹配和检查条件,则不会出现错误
标签: scala apache-spark pattern-matching case-class