【问题标题】:Spark parallel processing of grouped dataSpark并行处理分组数据
【发布时间】:2016-08-12 15:43:51
【问题描述】:

最初,我有很多数据。但是使用 spark-SQL 尤其是 groupBy 可以将其缩减到可管理的大小。 (适合单个节点的 RAM)

我如何在所有组上(并行)执行功能(分布在我的节点之间)?

如何确保将单个组的数据收集到单个节点?例如。我可能想使用local matrix 进行计算,但不想遇到有关数据局部性的错误。

【问题讨论】:

    标签: apache-spark apache-spark-sql apache-spark-mllib scala-breeze


    【解决方案1】:

    假设你有 x 没有。执行程序(在您的情况下,每个节点可能有 1 个执行程序)。并且您希望以这样的方式对密钥上的数据进行分区,使每个密钥都落入一个独特的存储桶中,这将是一个完美的分区器。没有通用的方法这样做,但如果有一些特定于您的数据的固有分布/逻辑,则有可能实现这一点。

    我处理过一个特定的案例,我发现 Spark 的内置哈希分区器在均匀分配密钥方面做得不好。所以我使用 Guava 编写了一个自定义分区器,如下所示:
      class FooPartitioner(partitions: Int) extends org.apache.spark.HashPartitioner(partitions: Int) {
        override def getPartition(key: Any): Int = {
          val hasherer = Hashing.murmur3_32().newHasher()
          Hashing.consistentHash(
            key match {
              case i: Int => hasherer.putInt(i).hash.asInt()
              case _ => key.hashCode
              },PARTITION_SIZE)
      }
     }
    

    然后我将此分区器实例添加为我正在使用的 combineBy 的参数,以便生成的 rdd 以这种方式分区。 这可以很好地将数据分配到 x 个存储桶,但我想不能保证每个存储桶只有 1 个密钥。

    如果您使用的是 Spark 1.6 并使用数据帧,您可以像这样定义一个 udf val hasher = udf((i:Int)=>Hashing.consistentHash(Hashing.murmur3_32().newHasher().putInt(i) .hash.asInt(),PARTITION_SIZE)) 然后做dataframe.repartition(hasher(keyThatYouAreUsing)) 希望这能提供一些入门提示。

    【讨论】:

    • 但我是否理解正确 dataFrame.groupBy("someKey") 会自动更改分区。如果我需要特殊分区,则需要应用自定义分区器。并且要并行计算一个函数,应该使用 UDF 吗?
    • 是的; spark dataframe 很可能会在 groupBy 之后使用哈希分区器进行重新分区。在 rdd api 中,您会注意到 groupBy 可以将分区器作为参数。理想情况下,如果我在 GroupBy 之后进行重新分区,则我的重新分区可以确保所有groupBy 的键 reqd 存在于同一个分区中,那么它不应该进行洗牌;但我不确定这是否可以依靠。 Spark 中固有的;一个函数将并行应用;所以是的,UDF 是并行计算的,您可以在 .select 或 .map 等中进行的任何其他转换也是如此
    • 谢谢。我需要进一步研究,但这听起来是一个很好的起点。
    【解决方案2】:

    我从Efficient UD(A)Fs with PySpark 找到了解决方案 这个博客

    1. mapPartitions 来分割数据;
    2. udaf 将 spark 数据帧转换为 pandas 数据帧;
    3. 在 udaf 中执行数据 etl 逻辑并返回 pandas 数据帧;
    4. udaf 会将 pandas 数据帧转换为 spark 数据帧;
    5. toDF() 合并结果 spark 数据帧并像 SaveAsTable 一样持久化;

    python df = df.repartition('guestid').rdd.mapPartitions(udf_calc).toDF()

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-02-05
    • 2016-08-14
    • 2020-05-02
    • 2017-05-17
    • 1970-01-01
    • 2018-06-06
    • 2017-12-18
    相关资源
    最近更新 更多