【问题标题】:Creating dataset based on different case classes [duplicate]根据不同的案例类创建数据集[重复]
【发布时间】: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:27​​5) 在 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


【解决方案1】:

试试:

trait AustraliaFile extends Serializable

case class Australiafile1(sectionName: String, profitCentre: String, valueAgainst: String, Status: String) extends AustraliaFile

case class Australiafile2(sectionName: String, profitCentre: String) extends AustraliaFile

您的类不是Serializable,但Spark 只能编写可序列化对象。此外,基于共同祖先的相关类始终是一个好主意,这样您就可以将 RDD 声明为 RDD[AustraliaFile] 而不是 RDD[Any]

另外,你的类匹配逻辑可以简化为

def mapper(line: String, recordLayoutClassToBeUsed: String) = {
  val fields = line.split(",")
  recordLayoutClassToBeUsed match {
     case ("Australiafile1") => Australiafile1(fields(0), fields(1), fields(2), fields(3))
    case ("Australiafile2") => Australiafile2(fields(0), fields(1))
  }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-07-21
    • 2020-04-24
    • 2019-11-08
    • 1970-01-01
    • 1970-01-01
    • 2019-10-25
    • 2016-09-15
    相关资源
    最近更新 更多