【问题标题】:Hive partitions, Spark partitions and joins in Spark - how they relateSpark 中的 Hive 分区、Spark 分区和连接 - 它们之间的关系
【发布时间】: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, ...)

我想在dateid 字段上加入table1table2

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


    【解决方案1】:

    一般回答,

    Spark 分区 - 大型分布式数据集的(逻辑)块。 Spark 为单个分区生成单个任务,该任务将在执行程序 JVM 中运行。

    Hive Partitions 是一种通过根据分区键(列)将表划分为不同部分来将表组织成分区的方法。分区使访问数据更加简单明了。

    可以调整的配置很少 -

    spark.sql.files.maxPartitionBytes - 读取文件时打包到单个分区的最大字节数(默认128MB)

    spark.sql.files.openCostInBytes - 打开文件的估计成本,以可以同时扫描的字节数来衡量。这用于将多个文件放入分区中。最好高估,那么小文件的分区会比大文件的分区快(先调度)。 (默认 4 MB)

    spark.sql.shuffle.partitions - 配置在为联接或聚合打乱数据时要使用的分区数。 (默认为 200)

    【讨论】:

      【解决方案2】:

      那么结果数据集将分别有多少个分区?分区等于 S3 中的对象数?

      无法回答您提供的信息。 in latest versions 的分区数量主要取决于 spark.sql.files.maxPartitionByte,尽管其他因素也可以发挥一些作用。

      Spark 是否能够利用连接的字段之一是 Hive 表中的分区键这一事实来优化连接?

      目前还没有 (Spark 2.3.0),但是 Spark 可以利用分桶 (DISTRIBUTE BY) 来优化连接。见How to define partitioning of DataFrame?。一旦 Data Source API v2 稳定下来,这可能会在未来发生变化。

      现在假设我使用的是 RDD (...) 同样,Spark 是否能够利用 Hive 表中的分区键在连接中使用的事实?

      一点也不。即使数据是分桶的 RDD 转换并且functional Dataset transformations 是黑盒子。无法应用任何优化并在此处应用。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2019-10-27
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-08-16
        • 2015-05-05
        相关资源
        最近更新 更多