【问题标题】:How to ignore bad avro files while reading from a folder into RDD从文件夹读取到 RDD 时如何忽略错误的 avro 文件
【发布时间】:2018-06-26 21:23:25
【问题描述】:

我正在从一个文件夹中读取一组 avro 文件,并且程序错误并带有错误消息。 //格式化不正确。

df =sqlContext.read.format("com.databricks.spark.avro").load("/data/hadoop20180516/22/abc*.avro").count()
[Stage 2:==================================================>(27818 + 4) / 28318]18/06/14 10:53:44 ERROR Executor: Exception in task 27900.0 in stage 2.0 (TID 27905)

java.io.IOException: 不是 Avro 数据文件

文件夹有 30K+ 个文件,其中一个文件可能已损坏。 我想忽略坏文件并继续加载文件的其余部分。s

我尝试使用 .option 命令
.option("badRecordsPath", "/tmp/badRecordsPath") 没有用。

有什么建议吗?

【问题讨论】:

    标签: pyspark


    【解决方案1】:

    我对python的了解不够,无法给你一个好的代码示例,但我在Scala中解决了这个问题,所以你可以试试:

    使用

    读取文件夹内的所有文件
    val paths = sparkContext.wholeTextFiles(folderPath).collect { case x: (String, String) => x._1 }.collect()
    

    这里我使用部分函数只获取键(文件路径),并再次收集以遍历字符串数组,而不是字符串的 RDD

    将每个文件加载为 DataFrame 并跳过失败的文件

    val filteredDFs = files.map { path =>
          Try(sparkSession
            .read
            .format(format)
            .options(options)
            .load(path)).toOption}.filter(_.isDefined).map(_.get)
    

    最后使用 union 创建一个包含所有先前 DF 的 DataFrame

    val finalDF = filteredDfs.reduce((df1, df2) => df1.union(df2))
    

    【讨论】:

    • 我从上面的代码中了解到的是,您首先获取文件路径,然后使用 try 块将文件加载到数据帧中。如果文件损坏,它将无法加载..我不确定pyspark中是否有try命令..让我来探索一下。
    • 是的,它就是这么做的。在地图中,在python中会是这样的: try: sparkSession.read.format(format).options(options).load(path) except: someOtherType 然后尝试过滤掉 someOtherType,所以最后你会有一个 DataFrame 的集合,然后你对它们进行联合。
    猜你喜欢
    • 2016-12-13
    • 1970-01-01
    • 1970-01-01
    • 2017-11-04
    • 2019-12-20
    • 2010-12-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多