【问题标题】:How to auto calculate numRepartition while using spark dataframe write如何在使用火花数据帧写入时自动计算 numRepartition
【发布时间】:2018-08-13 03:02:26
【问题描述】:

当我尝试将数据帧写入 Hive Parquet 分区表时

df.write.partitionBy("key").mode("append").format("hive").saveAsTable("db.table")

它会在HDFS中创建很多块,每个块只有很小的数据。

我了解它是如何进行的,因为每个 spark 子任务都会创建一个块,然后向它写入数据。

我也明白,块数会提高 Hadoop 性能,但达到阈值后也会降低性能。

如果我想自动设置 numPartition,有人有好主意吗?

numPartition = ??? // auto calc basing on df size or something
df.repartition("numPartition").write
  .partitionBy("key")
  .format("hive")
  .saveAsTable("db.table")

【问题讨论】:

  • 这几乎就像你问过How to master Apache-Spark?。选择正确的并行度是充分利用Spark 的全部功能的关键。 Here 是一个好的开始;它归结为您正在处理的数据量:列数、列类型、行数等。达到numPartitions 的指标需要时间和精力(命中和试验)。我从行数和数据大小 (GB) 开始预测它,然后从那里对其进行微调
  • 您的博客很棒,我确实在 df 转换方面实施了您的一些最佳实践。您的方法看起来不错,但就我而言,我有大量的离线数据管道,如果我忽略“重新分区”部分并在之后对其进行优化,这是一个不错的选择吗?
  • @Eric Yiwei Liu,该博客来自Umberto Griffo。如果你想一步一步地做事情,它是完全可以接受的,即。现在跳过repartition,稍后再访问。事实上,恕我直言,当您处理像Spark 这样的复杂框架时,最好采用这条路线:快速构建一个初步解决方案,然后从那里逐步改进以改进它。回想一下过早的优化是万恶之源
  • @y2k-shubham 非常肯定,压缩 CDH 警告到目前为止运行良好,我认为这对于未来的 spark 总是一个很好的功能,让我们更加专注于开发。

标签: apache-spark hadoop hive


【解决方案1】:

首先,当您已经在使用partitionBy(key) 时,为什么还要进行额外的重新分区步骤 - 您的数据将根据密钥进行分区。

通常,您可以按列值重新分区,这是一种常见的情况,有助于 reduceByKey、基于列值过滤等操作。例如,

val birthYears = List(
  (2000, "name1"),
  (2000, "name2"),
  (2001, "name3"),
  (2000, "name4"),
  (2001, "name5")
)
val df = birthYears.toDF("year", "name")

df.repartition($"year") 

【讨论】:

    【解决方案2】:

    默认情况下,spark 将为随机操作创建 200 个分区。因此,200 个文件/块(如果文件大小较小)将写入 HDFS。

    使用以下配置,根据您在 Spark 中的数据配置 shuffle 后要创建的分区数:

    spark.conf.set("spark.sql.shuffle.partitions", <Number of paritions>)
    

    例如:spark.conf.set("spark.sql.shuffle.partitions", "5"),因此 Spark 将创建 5 个分区并将 5 个文件写入 HDFS。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-03-30
      • 2020-08-09
      • 2018-02-15
      • 2021-10-15
      • 2018-02-12
      • 2018-11-27
      • 1970-01-01
      • 2019-04-12
      相关资源
      最近更新 更多