【问题标题】:Catch Exceptions that are thrown on map function in Spark捕获 Spark 中 map 函数引发的异常
【发布时间】:2017-10-03 00:18:19
【问题描述】:

我正在读取一个包含一些损坏数据的文件。我的文件如下所示:

30149;E;LDI0775100  350000003221374461
30153;168034601 350000003486635135

第二行是应该的样子。第一行在第一列中有一些额外的字符。所以我想捕获由于数据损坏而引发的任何异常。不仅是上面的示例。下面是我将文件加载到我的 RDD 并尝试在 map 函数中捕获异常的代码。

  val rawCustfile = sc.textFile("/tmp/test_custmap")

  case class Row1(file_id:Int,mk_cust_id: String,ind_id:Long)


   val cleanedcustmap = rawCustfile.map(x => x.replaceAll(";", 
 "\t").split("\t")).map(x => Try{
     Row1(x(0).toInt, x(1), x(2).toLong)}match {
      case Success(map) => Right(map)
      case Failure(e) => Left(e)
    }) 

    //get the good columns
   val good_rows=cleanedcustmap.filter(_.isRight)

    //get the errors
    val error=cleanedcustmap.filter(_.isLeft)

     good_rows.collect().foreach(println)

   error.collect().foreach(println)

   val df = sqlContext.createDataFrame(cleanedcustmap.filter(_.isRight))

这个 good_rows.collect().foreach(println) 打印:

 Right(Row1(30153,168034601,350000003486635135))

这个 error.collect().foreach(println) 打印:

 Left(java.lang.NumberFormatException: For input string: "LDI0775100")

在我尝试将我的 rdd 转换为 DataFrame 之前,一切正常。我得到以下异常:

 Name: scala.MatchError
 Message: Product with Serializable with 
 scala.util.Either[scala.Throwable,Row1] (of class scala.reflect.internal.Types$RefinedType0)
    StackTrace: org.apache.spark.sql.catalyst.ScalaReflection$class.schemaFor(ScalaReflection.scala:676)
    org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:30)
    org.apache.spark.sql.catalyst.ScalaReflection$class.schemaFor(ScalaReflection.scala:630)
    org.apache.spark.sql.catalyst.ScalaReflection$.schemaFor(ScalaReflection.scala:30)
    org.apache.spark.sql.SQLContext.createDataFrame(SQLContext.scala:414)
    $line79.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC.<init>(<console>:44)
    $line79.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC.<init>(<console>:49)
    $line79.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC.<init>(<console>:51)
    $line79.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC.<init>(<console>:53)
    $line79.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC.<init>(<console>:55)
    $line79.$read$$iwC$$iwC$$iwC$$iwC$$iwC.<init>(<console>:57)
    $line79.$read$$iwC$$iwC$$iwC$$iwC.<init>(<console>:59)
    $line79.$read$$iwC$$iwC$$iwC.<init>(<console>:61)
    $line79.$read$$iwC$$iwC.<init>(<console>:63)
    $line79.$read$$iwC.<init>(<console>:65)
    $line79.$read.<init>(<console>:67)
    $line79.$read$.<init>(<console>:71)
    $line79.$read$.<clinit>(<console>)
    $line79.$eval$.<init>(<console>:7)
    $line79.$eval$.<clinit>(<console>)
    $line79.$eval.$print(<console>)
    sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    java.lang.reflect.Method.invoke(Method.java:497)
    org.apache.spark.repl.SparkIMain$ReadEvalPrint.call(SparkIMain.scala:1065)
    org.apache.spark.repl.SparkIMain$Request.loadAndRun(SparkIMain.scala:1346)
    org.apache.spark.repl.SparkIMain.loadAndRunReq$1(SparkIMain.scala:840)
    org.apache.spark.repl.SparkIMain.interpret(SparkIMain.scala:871)
    org.apache.spark.repl.SparkIMain.interpret(SparkIMain.scala:819)
    org.apache.toree.kernel.interpreter.scala.ScalaInterpreter$$anonfun$interpretAddTask$1$$anonfun$apply$3.apply(ScalaInterpreter.scala:361)
    org.apache.toree.kernel.interpreter.scala.ScalaInterpreter$$anonfun$interpretAddTask$1$$anonfun$apply$3.apply(ScalaInterpreter.scala:356)
    org.apache.toree.global.StreamState$.withStreams(StreamState.scala:81)
    org.apache.toree.kernel.interpreter.scala.ScalaInterpreter$$anonfun$interpretAddTask$1.apply(ScalaInterpreter.scala:355)
    org.apache.toree.kernel.interpreter.scala.ScalaInterpreter$$anonfun$interpretAddTask$1.apply(ScalaInterpreter.scala:355)
    org.apache.toree.utils.TaskManager$$anonfun$add$2$$anon$1.run(TaskManager.scala:140)
    java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
    java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
    java.lang.Thread.run(Thread.java:745)

我的第一个问题是,我是否以正确的方式捕捉异常?我想得到错误,因为我想打印它们。有没有更好的办法?。我的第二个问题是将我的 RDD 转换为 DataFrama 时我做错了什么

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    我是否以正确的方式捕捉异常

    嗯,差不多。不确定TryEither 的映射是否重要(当左侧类型为Throwable 时,您可以将Try 视为Either 的特化...)但两者能行得通。您可能想要解决的另一个问题是使用 replaceAll - 它从第一个参数(在这种情况下不需要)生成 regex,因此比 replace 慢,见Difference between String replace() and replaceAll()

    将我的 RDD 转换为 DataFrame 时我做错了什么

    DataFrames 仅支持一组有限的“标准”类型

    • 基元(例如整数、长整数、字符串...)
    • 数组/映射(其他支持的类型)
    • 产品(案例类或其他受支持类型的元组)

    Either 不是这些类型之一(Try 也不是),因此它不能在 DataFrame 中使用

    您可以通过以下方式解决:

    1. 仅使用 DataFrame 中的“成功”记录(除了日志记录,这些错误还有什么用处?无论哪种方式,它们都应该单独处理):

      // this would work:
      val df = spark.createDataFrame(good_rows.map(_.right.get))
      
    2. Either 转换为支持的类型,例如Tuple2(Row1, String),其中字符串是错误消息,元组中的值之一是null(左侧为错误记录为空,右侧为成功记录为空)

    【讨论】:

    • 非常感谢您的回复。我很抱歉,因为我在创建 DataFrame 时粘贴了错误的行。我确实尝试了你的第一个建议,但我得到了上面提到的例外。如果造成任何混乱,我深表歉意。我已经更新了我的问题
    • 过滤 isReight 是不够的 - 它只会传递具有正确值的 Either 实例,但这些仍然是 Either 的实例。看起来已关闭,我建议的解决方案 - 它还使用 map(_.right.get)映射这些记录,将 Either 实例“解包”到底层 Row1 值中。
    • 您甚至可以使用good_rows = cleanedcustmap.collect{case Right(r) =&gt; r} 来避免使用不安全的right.get 方法(不是在这里,但一般来说,如果您忘记在某个时候过滤掉Lefts)。
    猜你喜欢
    • 1970-01-01
    • 2019-09-13
    • 1970-01-01
    • 1970-01-01
    • 2019-07-17
    • 1970-01-01
    • 1970-01-01
    • 2012-06-05
    • 1970-01-01
    相关资源
    最近更新 更多