【问题标题】:Exception is not caught by the try catch blocktry catch 块没有捕获到异常
【发布时间】:2019-07-12 21:30:45
【问题描述】:

我将 DStream 保存到 Cassandra。 Cassandra 中有一个带有map<text, text> 数据类型的列。 Cassandra 不支持 Map 中的 null 值,但流中可能出现空值。

如果出现问题,我已添加 try catch,但尽管如此,程序还是停止了,我在日志中没有看到错误消息:

   try {
      cassandraStream.saveToCassandra("table", "keyspace")
    } catch {
      case e: Exception => log.error("Error in saving data in Cassandra" + e.getMessage, e)
    }

例外

Caused by: java.lang.NullPointerException: Map values cannot be null
    at com.datastax.driver.core.TypeCodec$AbstractMapCodec.serialize(TypeCodec.java:2026)
    at com.datastax.driver.core.TypeCodec$AbstractMapCodec.serialize(TypeCodec.java:1909)
    at com.datastax.driver.core.AbstractData.set(AbstractData.java:530)
    at com.datastax.driver.core.AbstractData.set(AbstractData.java:536)
    at com.datastax.driver.core.BoundStatement.set(BoundStatement.java:870)
    at com.datastax.spark.connector.writer.BoundStatementBuilder.com$datastax$spark$connector$writer$BoundStatementBuilder$$bindColumnUnset(BoundStatementBuilder.scala:73)
    at com.datastax.spark.connector.writer.BoundStatementBuilder$$anonfun$6.apply(BoundStatementBuilder.scala:84)
    at com.datastax.spark.connector.writer.BoundStatementBuilder$$anonfun$6.apply(BoundStatementBuilder.scala:84)
    at com.datastax.spark.connector.writer.BoundStatementBuilder$$anonfun$bind$1.apply$mcVI$sp(BoundStatementBuilder.scala:106)
    at scala.collection.immutable.Range.foreach$mVc$sp(Range.scala:160)
    at com.datastax.spark.connector.writer.BoundStatementBuilder.bind(BoundStatementBuilder.scala:101)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.next(GroupingBatchBuilder.scala:106)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.next(GroupingBatchBuilder.scala:31)
    at scala.collection.Iterator$class.foreach(Iterator.scala:893)
    at com.datastax.spark.connector.writer.GroupingBatchBuilder.foreach(GroupingBatchBuilder.scala:31)
    at com.datastax.spark.connector.writer.TableWriter$$anonfun$writeInternal$1.apply(TableWriter.scala:233)
    at com.datastax.spark.connector.writer.TableWriter$$anonfun$writeInternal$1.apply(TableWriter.scala:210)
    at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$withSessionDo$1.apply(CassandraConnector.scala:112)
    at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$withSessionDo$1.apply(CassandraConnector.scala:111)
    at com.datastax.spark.connector.cql.CassandraConnector.closeResourceAfterUse(CassandraConnector.scala:145)
    at com.datastax.spark.connector.cql.CassandraConnector.withSessionDo(CassandraConnector.scala:111)
    at com.datastax.spark.connector.writer.TableWriter.writeInternal(TableWriter.scala:210)
    at com.datastax.spark.connector.writer.TableWriter.insert(TableWriter.scala:197)
    at com.datastax.spark.connector.writer.TableWriter.write(TableWriter.scala:183)
    at com.datastax.spark.connector.streaming.DStreamFunctions$$anonfun$saveToCassandra$1$$anonfun$apply$1.apply(DStreamFunctions.scala:54)
    at com.datastax.spark.connector.streaming.DStreamFunctions$$anonfun$saveToCassandra$1$$anonfun$apply$1.apply(DStreamFunctions.scala:54)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
    at org.apache.spark.scheduler.Task.run(Task.scala:109)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:345)
    ... 3 more

我想知道为什么程序会停止,尽管有 try/catch 块。为什么没有捕获到异常?

【问题讨论】:

    标签: scala apache-spark exception-handling try-catch spark-streaming


    【解决方案1】:

    问题是你没有捕捉到你认为你会捕捉到的异常。您拥有的代码将捕获驱动程序异常,事实上,像这样构造的代码会做到这一点。

    但这并不意味着

    程序不应停止。

    虽然驱动程序故障(可能是致命的执行程序故障的结果)得到了控制,并且驱动程序可以正常退出,但流本身已经消失了。因此您的代码退出,因为没有更多的流可以运行。

    如果有问题的代码在您的控制之下,则应将异常处理委托给该任务,但如果是第 3 方代码,则没有这样的选项。

    相反,您应该验证您的数据并删除有问题的记录,然后再将它们传递给saveToCassandra

    【讨论】:

    • 如果cassandraStreamRDD,那是真的。但是,如果它是DStream,则异常抛出与try / catch 保护的代码没有直接关系。
    【解决方案2】:

    要了解失败的根源,您必须承认DStreamFunctions.saveToCassandra 与一般的DStream 输出操作相同,不是严格意义上的操作。在实践中it just invokes foreachRDD

    dstream.foreachRDD(rdd => rdd.sparkContext.runJob(rdd, writer.write _))
    

    which in turn:

    对这个 DStream 中的每个 RDD 应用一个函数。这是一个输出运算符,因此“this”DStream 将被注册为输出流并因此实现。

    区别很微妙,但很重要 - 操作已注册,但实际执行发生在不同的上下文中,在稍后的时间点。

    这意味着在您调用 saveToCassandra 时没有捕获到运行时故障。

    如前所述,tryTry 将包含驱动程序异常,如果直接应用于操作。因此,例如,您将 saveToCassandra 重新实现为

    dstream.foreachRDD(rdd => try { 
      rdd.sparkContext.runJob(rdd, writer.write _) 
    } catch {
      case e: Exception => log.error("Error in saving data in Cassandra" + e. getMessage, e)
    })
    

    流应该能够继续,尽管当前批次将完全或部分丢失。

    需要注意的是,这与捕获原始异常不同,它将被抛出,未被捕获并且在日志中可见。要从源头捕获问题,您必须直接在编写器中应用 try / catch 块,这显然不是您执行代码时的选择,您无法控制它。

    带走消息(已在此线程中说明) - 确保清理您的数据以避免已知的故障来源。

    【讨论】:

    • 我正在使用用于 DStream 的 saveToCassandra。该函数不支持rdd
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-06-26
    • 2019-11-11
    • 1970-01-01
    • 2010-09-07
    • 1970-01-01
    相关资源
    最近更新 更多