如何让spark忽略丢失的输入文件?

jbr*_*own 11 hadoop apache-spark

我想在包含avro文件的一些生成的S3路径上运行一个spark作业(spark v1.5.1).我正在加载它们:

val avros = paths.map(p => sqlContext.read.avro(p))
Run Code Online (Sandbox Code Playgroud)

但是有些路径不存在.我如何能够忽略那些空路径?以前我已经使用过这个答案,但我不确定如何在新的数据帧API中使用它.

注意:我理想地寻找一种类似于链接答案的方法,它只是使输入路径可选.我并不特别想要在S3中明确检查路径是否存在(因为这很麻烦并且可能使开发变得尴尬),但我想这是我的后备,如果现在没有干净的方法来实现它.

mat*_*its 12

我会使用scala Try类型来处理读取avro文件目录时失败的可能性.使用'Try',我们可以在代码中明确失败的可能性,并以功能方式处理它:

object Main extends App {

  import scala.util.{Success, Try}
  import org.apache.spark.{SparkConf, SparkContext}
  import com.databricks.spark.avro._

  val sc = new SparkContext(new SparkConf().setMaster("local[*]").setAppName("example"))
  val sqlContext = new org.apache.spark.sql.SQLContext(sc)

  //the first path exists, the second one doesn't
  val paths = List("/data/1", "/data/2")

  //Wrap the attempt to read the paths in a Try, then use collect to filter
  //and map with a single partial function.
  val avros =
    paths
      .map(p => Try(sqlContext.read.avro(p)))
      .collect{
        case Success(df) => df
      }
  //Do whatever you want with your list of dataframes
  avros.foreach{ df =>
    println(df.collect())
  }
  sc.stop()
}
Run Code Online (Sandbox Code Playgroud)

  • 当在RDD上调用`collect()`时就是这样.我第一次调用`collect(...)`,包含部分函数,​​它在RDD列表中,它是List上的collect函数,而不是任何RDD.这相当于做一个`map`和`filter`.我最后在`foreach`的末尾再次使用`collect()`,但这只是操作RDD列表的一个例子,我不希望你在自己的应用程序中做什么,但是我需要一个简单的结局才能看出该方法是否正常工作. (2认同)