【发布时间】:2017-12-17 10:53:00
【问题描述】:
在我们的应用程序中,我们的大部分代码只是在DataFrame上应用filter、group by和aggregate操作并将DF保存到Cassandra数据库。
像下面的代码一样,我们有几个方法可以对不同数量的字段执行相同类型的操作[filter, group by, join, agg],并返回一个 DF,并将其保存到 Cassandra 表中。
示例代码是:
val filteredDF = df.filter(col("hour") <= LocalDataTime.now().getHour())
.groupBy("country")
.agg(sum(col("volume")) as "pmtVolume")
saveToCassandra(df)
def saveToCassandra(df: DataFrame) {
try {
df.write.format("org.apache.spark.sql.cassandra")
.options(Map("Table" -> "tableName", "keyspace" -> keyspace)
.mode("append").save()
}
catch {
case e: Throwable => log.error(e)
}
}
由于我通过将 DF 保存到 Cassandra 来调用操作,我希望我只需要按照 this 线程处理该行上的异常。
如果我得到任何异常,我可以默认在 Spark 详细日志中看到异常。
我是否必须真正包围过滤器,使用Try或try , catch?按代码分组
我在 Spark SQL DataFrame API 示例中没有看到任何带有异常处理的示例。
如何在saveToCassandra 方法上使用Try?它返回Unit
【问题讨论】:
标签: scala exception-handling apache-spark-sql