【问题标题】:Spark SQL DataFrame - Exception handlingSpark SQL DataFrame - 异常处理
【发布时间】:2017-12-17 10:53:00
【问题描述】:

在我们的应用程序中,我们的大部分代码只是在DataFrame上应用filtergroup byaggregate操作并将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 详细日志中看到异常。

我是否必须真正包围过滤器,使用Trytry , catch?按代码分组

我在 Spark SQL DataFrame API 示例中没有看到任何带有异常处理的示例。

如何在saveToCassandra 方法上使用Try?它返回Unit

【问题讨论】:

    标签: scala exception-handling apache-spark-sql


    【解决方案1】:

    在 try catch 中包装惰性 DAG 是没有意义的。
    您需要将 lambda 函数包装在 Try() 中。
    不幸的是,AFAIK 无法在 DataFrames 中进行行级异常处理。

    您可以使用 RDD 或 DataSet,如下面的帖子中所述 spache spark exception handling

    【讨论】:

      【解决方案2】:

      您实际上不需要用Trytrycatch 包围filtergroup by 代码。由于所有这些操作都是转换,因此在对它们执行 action 之前,它们不会被执行,例如 saveToCassandra 在您的情况下。

      但是,如果在过滤分组聚合数据帧时发生错误,saveToCassandra 函数中的 catch 子句将记录它正在那里执行操作。

      【讨论】:

        猜你喜欢
        • 2020-03-08
        • 2016-05-26
        • 2015-12-16
        • 2017-12-20
        • 2016-10-10
        • 1970-01-01
        • 1970-01-01
        • 2017-01-30
        • 1970-01-01
        相关资源
        最近更新 更多