【发布时间】: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,我有多个输出表,base、A、B 等在基于作业时间戳的给定路径中。
我想left join 他们都基于时间戳和主目录,在本例中为foo。这意味着将每个输出表 base、A、B 等读入新的单独输入表中,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 序列的正确方向吗?甚至可能值得将读取视为惰性或顺序读取,因此仅在应用连接时读取A 或B 表以减少内存需求。
注意:目录结构不是最终的,这意味着如果适合解决方案,它可以更改。
【问题讨论】:
-
看起来像 hive 分区结构,您正在使用 orc 日期分区文件。为什么不能将这些映射到 hive 并为每个日期使用
hiveContext.sql然后加入 -
我们没有运行 Hive,只运行 Spark 独立
标签: scala apache-spark amazon-s3 apache-spark-sql orc