【发布时间】:2018-10-05 13:12:11
【问题描述】:
试图了解 Hive 分区与 Spark 分区之间的关系,最终提出了一个关于连接的问题。
我有 2 个外部 Hive 表;均由 S3 存储桶支持并由 date 分区;所以在每个桶中都有名称格式为date=<yyyy-MM-dd>/<filename>的键。
问题 1:
如果我将这些数据读入 Spark:
val table1 = spark.table("table1").as[Table1Row]
val table2 = spark.table("table2").as[Table2Row]
那么结果数据集将分别有多少个分区?分区等于 S3 中的对象数?
问题 2:
假设这两种行类型具有以下架构:
Table1Row(date: Date, id: String, ...)
Table2Row(date: Date, id: String, ...)
我想在date 和id 字段上加入table1 和table2:
table1.joinWith(table2,
table1("date") === table2("date") &&
table1("id") === table2("id")
)
Spark 是否能够利用连接的字段之一是 Hive 表中的分区键这一事实来优化连接?如果是的话怎么办?
问题 3:
现在假设我改用RDDs:
val rdd1 = table1.rdd
val rdd2 = table2.rdd
AFAIK,使用RDD API 的连接语法类似于:
rdd1.map(row1 => ((row1.date, row1.id), row1))
.join(rdd2.map(row2 => ((row2.date, row2.id), row2))))
同样,Spark 是否能够利用 Hive 表中的分区键在连接中使用的事实?
【问题讨论】:
标签: apache-spark hive apache-spark-sql apache-spark-dataset