【问题标题】:scala.MatchError: [abc,cde,null,3] (of class org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema) in Spark JSON with missing fieldsscala.MatchError: [abc,cde,null,3] (of class org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema) 在 Spark JSON 中缺少字段
【发布时间】:2018-03-05 05:34:26
【问题描述】:

我有 JSON 输入文件:

{"a": "abc", "b": "bcd", "d": 3},
{"a": "ezx", "b": "hdg", "c": "ssa"},
...

每个对象的某些字段丢失,而不是放置 null 值。

在使用 Scala 的 Apache Spark 中:

import SparkCommons.sparkSession.implicits._

private val inputJsonPath: String = "resources/input/input.json"

private val schema = StructType(Array(
  StructField("a", StringType, nullable = false),
  StructField("b", StringType, nullable = false),
  StructField("c", StringType, nullable = true),
  StructField("d", DoubleType, nullable = true)
))

private val inputDF: DataFrame = SparkCommons.sparkSession
  .read
  .schema(schema)
  .json(inputJsonPath)
  .cache()

inputDF.printSchema()

val dataRdd = inputDF.rdd
.map {
  case Row(a: String, b: String, c: String, d: Double) =>
    MyCaseClass(a, b, c, d)
}

val dataMap = dataRdd.collectAsMap()

MyCaseClass 代码:

case class MyCaseClass(
              a: String,
              b: String,
              c: String = null,
              d: Double = Predef.Double2double(null)
)

我得到以下模式作为输出:

root
 |-- a: string (nullable = true)
 |-- b: string (nullable = true)
 |-- c: string (nullable = true)
 |-- d: double (nullable = true)

程序可以编译,但在运行时 Spark 完成工作后,我会收到以下异常:

[error] - org.apache.spark.executor.Executor - Exception in task 3.0 in stage 4.0 (TID 21)
scala.MatchError: [abc,bcd,null,3] (of class org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema)
at com.matteoguarnerio.spark.SparkOperations$$anonfun$1.apply(SparkOperations.scala:62) ~[classes/:na]
at com.matteoguarnerio.spark.SparkOperations$$anonfun$1.apply(SparkOperations.scala:62) ~[classes/:na]
at scala.collection.Iterator$$anon$11.next(Iterator.scala:410) ~[scala-library-2.11.11.jar:na]
at scala.collection.Iterator$$anon$11.next(Iterator.scala:410) ~[scala-library-2.11.11.jar:na]
at scala.collection.Iterator$$anon$11.next(Iterator.scala:410) ~[scala-library-2.11.11.jar:na]
at org.apache.spark.util.random.SamplingUtils$.reservoirSampleAndCount(SamplingUtils.scala:42) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.RangePartitioner$$anonfun$9.apply(Partitioner.scala:261) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.RangePartitioner$$anonfun$9.apply(Partitioner.scala:259) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsWithIndex$1$$anonfun$apply$25.apply(RDD.scala:820) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsWithIndex$1$$anonfun$apply$25.apply(RDD.scala:820) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:319) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.rdd.RDD.iterator(RDD.scala:283) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:70) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.scheduler.Task.run(Task.scala:86) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:274) ~[spark-core_2.11-2.0.2.jar:2.0.2]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [na:1.8.0_144]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [na:1.8.0_144]
at java.lang.Thread.run(Thread.java:748) [na:1.8.0_144]

Spark 版本:2.0.2

Scala 版本:2.11.11

  • 即使在RDD匹配和创建对象中某些字段为null或缺失,如何解决此异常并进行迭代?
  • 为什么架构,即使我在某些字段上明确定义不可为空和可空,一切都可以为空?

更新

我刚刚在dataRdd 上使用了一种解决方法来避免这个问题:

private val dataRdd = inputDF.rdd
.map {
  case r: GenericRowWithSchema => {
      val a = r.getAs("a").asInstanceOf[String]
      val b = r.getAs("b").asInstanceOf[String]

      var c: Option[String] = None
      var d: Option[Double] = None

      try {
        c = if (r.isNullAt(r.fieldIndex("c"))) None: Option[String] else Some(r.getAs("c").asInstanceOf[String])
        d = if (r.isNullAt(r.fieldIndex("d"))) None: Option[Double] else Some(r.getAs("d").asInstanceOf[Double])
      } catch {
        case _: Throwable => None
      }

      MyCaseClass(a, b, c, d)
  }
}

并以这种方式更改MyCaseClass

case class MyCaseClass(
              a: String,
              b: String,
              c: Option[String],
              d: Option[Double]
)

【问题讨论】:

    标签: json scala apache-spark nullpointerexception pattern-matching


    【解决方案1】:

    问题在于input.json。应该是这样的:

    {"a": "abc", "b": "bcd", "d": 3},
    {"a": "ezx", "b": "hdg", "c": "ssa"},
    ...
    

    有了这个input.json,你的代码就可以正常工作了。

    【讨论】:

    • JSON 已经被引用为属性键。我错误地转储了文件,它已经是那种形式了,问题仍然存在。
    • 就像我在回答中提到的 - 您的代码运行良好!除非您的问题中缺少其他内容,否则没有问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-06-20
    • 1970-01-01
    • 2023-03-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-30
    相关资源
    最近更新 更多