【发布时间】: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