【问题标题】:Spark read multiple directories into multiple dataframesSpark将多个目录读入多个数据帧
【发布时间】:2017-06-23 02:48:21
【问题描述】:

我在 S3 上的目录结构如下所示:

foo
  |-base
     |-2017
        |-01
           |-04
              |-part1.orc, part2.orc ....
  |-A
     |-2017
        |-01
           |-04
              |-part1.orc, part2.orc ....
  |-B
     |-2017
        |-01
           |-04
              |-part1.orc, part2.orc ....

这意味着对于目录foo,我有多个输出表,baseAB 等在基于作业时间戳的给定路径中。

我想left join 他们都基于时间戳和主目录,在本例中为foo。这意味着将每个输出表 baseAB 等读入新的单独输入表中,left join 可以在这些输入表中应用。全部以base 表为起点

类似这样的东西(不工作的代码!)

val dfs: Seq[DataFrame] = spark.read.orc("foo/*/2017/01/04/*")
val base: DataFrame = spark.read.orc("foo/base/2017/01/04/*")

val result = dfs.foldLeft(base)((l, r) => l.join(r, 'id, "left"))

有人可以为我指出如何获取 DataFrame 序列的正确方向吗?甚至可能值得将读取视为惰性或顺序读取,因此仅在应用连接时读取AB 表以减少内存需求。

注意:目录结构不是最终的,这意味着如果适合解决方案,它可以更改。

【问题讨论】:

  • 看起来像 hive 分区结构,您正在使用 orc 日期分区文件。为什么不能将这些映射到 hive 并为每个日期使用 hiveContext.sql 然后加入
  • 我们没有运行 Hive,只运行 Spark 独立

标签: scala apache-spark amazon-s3 apache-spark-sql orc


【解决方案1】:

据我了解,Spark 使用底层 Hadoop API 来读取数据文件。因此,继承的行为是将您指定的所有内容读入一个 RDD/DataFrame。

要实现你想要的,你可以先得到一个目录列表:

    import org.apache.hadoop.conf.Configuration
    import org.apache.hadoop.fs.{ FileSystem, Path }

    val path = "foo/"

    val hadoopConf = new Configuration()
    val fs = FileSystem.get(hadoopConf)
    val paths: Array[String] = fs.listStatus(new Path(path)).
      filter(_.isDirectory).
      map(_.getPath.toString)

然后将它们加载到单独的数据帧中:

    val dfs: Array[DataFrame] = paths.
      map(path => spark.read.orc(path + "/2017/01/04/*"))

【讨论】:

    【解决方案2】:

    这是(我认为)您正在尝试做的事情的直接解决方案,不使用 Hive 等额外功能或内置分区功能:

    import spark.implicits._
    
    // load base
    val baseDF = spark.read.orc("foo/base/2017/01/04").as("base")
    
    // create or use existing Hadoop FileSystem - this should use the actual config and path
    val fs = FileSystem.get(new URI("."), new Configuration())
    
    // find all other subfolders under foo/
    val otherFolderPaths = fs.listStatus(new Path("foo/"), new PathFilter {
      override def accept(path: Path): Boolean = path.getName != "base"
    }).map(_.getPath)
    
    // use foldLeft to join all, using the DF aliases to find the right "id" column
    val result = otherFolderPaths.foldLeft(baseDF) { (df, path) =>
      df.join(spark.read.orc(s"$path/2017/01/04").as(path.getName), $"base.id" === $"${path.getName}.id" , "left") }
    

    【讨论】:

      猜你喜欢
      • 2023-03-25
      • 2013-11-30
      • 2012-06-28
      • 1970-01-01
      • 2020-05-01
      • 2016-10-03
      • 1970-01-01
      • 2018-09-26
      相关资源
      最近更新 更多