【问题标题】:Pre-cogrouping tables on HDFS and reading in Spark with zero shuffling在 HDFS 上预先组合表并在 Spark 中读取零混洗
【发布时间】:2019-01-31 01:30:46
【问题描述】:

上下文

我有两个表作为我的 spark 作业的一部分加入/联合分组,每次运行作业时都会产生很大的洗牌。我想通过存储一次共同分组的数据来分摊所有作业的成本,并将已经共同分组的数据用作我常规 Spark 运行的一部分以避免洗牌。

为了尝试实现这一点,我在 HDFS 中以 parquet 格式存储了一些数据。我正在使用 Parquet 重复字段来实现以下架构

(日期,[aRecords],[bRecords])

其中 [aRecords] 表示 aRecord 的数组。我还使用通常的 write.partitionBy($"date") 在 HDFS 上按日期对数据进行分区。

在这种情况下,aRecords 和 bRecords 似乎按日期有效地组合在一起。我可以执行如下操作:

case class CogroupedData(date: Date, aRecords: Array[Int], bRecords: Array[Int])

val cogroupedData = spark.read.parquet("path/to/data").as[CogroupedData]

//Dataset[(Date,Int)] where the Int in the two sides multiplied
val results = cogroupedData
    .flatMap(el => el.aRecords.zip(el.bRecords).map(pair => (el.date, pair._1 * pair._2)))

并获得我通过在两个单独的表上使用等效的 groupByKey 操作获得的结果,这些表分别是按日期键入的 aRecords 和 bRecords。

两者之间的区别在于,我避免了对已经同组的数据进行洗牌,同组的成本通过在 HDFS 上持久化来摊销。

问题

现在回答问题。从 cogrouped 数据集中,我想派生两个分组数据集,以便我可以使用标准 Spark SQL 运算符(如 cogroup、join 等)而不会产生随机播放。这似乎是可能的,因为第一个代码示例有效,但是当我加入/groupByKey/cogroup 等时,Spark 仍然坚持对数据进行散列/洗牌。

以下面的代码示例为例。我希望有一种方法可以在执行连接时运行以下内容而不会产生随机播放。

val cogroupedData = spark.read.parquet("path/to/data").as[CogroupedData]

val aRecords = cogroupedData
    .flatMap(cog => cog.aRecords.map(a => (cog.date,a)))
val bRecords = cogroupedData
    .flatMap(cog => cog.bRecords.map(b => (cog.date,b)))

val joined = aRecords.join(bRecords,Seq("date"))

查看文献,如果 cogroupedData 有一个已知的分区器,那么接下来的操作不应该引起洗牌,因为它们可以使用 RDD 已经分区的事实并保留分区器。

我认为我需要实现这一目标是获得一个带有已知分区器的 cogroupedData 数据集/rdd,而不会产生洗牌。

我已经尝试过的其他事情:

  • Hive 元数据 - 适用于简单连接,但仅优化初始连接而不是后续转换。 Hive 对 cogroup 也毫无帮助

有人有什么想法吗?

【问题讨论】:

    标签: scala apache-spark hadoop apache-spark-sql


    【解决方案1】:

    你在这里犯了两个错误。

    正确的做法是:

    • 使用分桶。

      val n: Int
      someDF.write.bucketBy(n, "date").saveAsTable("df")
      
    • 放弃功能性 API 以支持 SQL API:

      import org.apache.spark.sql.functions.explode
      
      val df = spark.table("df")
      
      val adf = df.select($"date", explode($"aRecords").alias("aRecords"))
      val bdf = df.select($"date", explode($"bRecords").alias("bRecords"))
      
      adf.join(bdf, Seq("date"))
      

    【讨论】:

      猜你喜欢
      • 2015-09-03
      • 1970-01-01
      • 1970-01-01
      • 2015-01-26
      • 1970-01-01
      • 2017-08-07
      • 1970-01-01
      • 1970-01-01
      • 2015-12-18
      相关资源
      最近更新 更多