【发布时间】:2019-07-25 17:25:40
【问题描述】:
假设我有一些数据都在同一个分区上(我之前在数据帧上执行了.coalesce(1))。我现在想对数据进行分组并对其执行聚合。如果我在数据帧上使用.groupBy,这些组会被放置到不同的节点上吗?
如果它是真的,我想避免这种情况,因为我想对组执行这些计算而不用太多洗牌。
【问题讨论】:
标签: scala apache-spark hadoop apache-spark-sql bigdata
假设我有一些数据都在同一个分区上(我之前在数据帧上执行了.coalesce(1))。我现在想对数据进行分组并对其执行聚合。如果我在数据帧上使用.groupBy,这些组会被放置到不同的节点上吗?
如果它是真的,我想避免这种情况,因为我想对组执行这些计算而不用太多洗牌。
【问题讨论】:
标签: scala apache-spark hadoop apache-spark-sql bigdata
首先,coalesce(1) 不能保证您的所有数据都在一个节点中,因此您必须使用repartition(1),这将强制将您的所有数据合并到一个节点中。 coalesce 只对同一个节点中的分区进行分组,所以如果你的数据分布在 5 个节点中(每个节点有多个分区),最后会保留 5 个分区。 repartition 强制洗牌,将所有数据移动到单个节点。
但是,如果您关心的是聚合中的分区数量,这取决于,如果聚合只是您所有数据的reduce,spark sql 将尝试首先在每个节点中减少,然后减少结果每个节点,一个例子将是一个计数。但是对于分桶聚合,例如计算具有 id 的元素的数量,spark 所做的是首先在每个节点中减少,然后将数据打乱到桶中,以确保每个节点的所有减少,对于相同的 id在同一个节点中,并再次减少它们。存储桶的数量由属性spark.sql.shuffle.partitions 配置,每个存储桶都将作为您的作业中的一个任务执行。请小心,因为将 spark.sql.shuffle.partitions 设置为 1 可能会使进程的其他部分(例如连接或大聚合)变慢,或者导致内存不足错误。
【讨论】:
repartition(1) 相当于将spark.sql.shuffle.partitions 设置为1?所以在repartition(1) 之后执行groupBy 时不会重新洗牌?
这取决于。默认情况下,分区数由 spark.sql.shuffle.partitions 定义。避免这种情况的一种方法是使用带有显式分区表达式的repartition 而不是coalesce:
val df = sparkSession.createDataFrame(
sparkContext.parallelize(Seq(Row(1, "a"), Row(1, "b"), Row(2, "c"))),
StructType(List(StructField("foo", IntegerType, true), StructField("bar", StringType, true))))
df.repartition(numPartitions = 1, $"foo").groupBy("foo").agg(count("*")).explain()
一般来说,可以使用 Spark Web UI 并在“阶段”选项卡上监控随机读取/写入指标。
【讨论】:
repartition 和 coalesce 这里到底有什么区别?