【发布时间】:2021-09-07 13:38:24
【问题描述】:
我正在用非常简单的逻辑进行数据清理。
val inputData= spark.read.parquet(inputDataPath)
val viewMiddleTable = sdk70000DF.where($"type" === "0").select($"field1", $"field2", $field3)
.groupBy($"field1", $"field2", $field3)
.agg(count(lit(1)))
从hdfs读取parquet数据,过滤,选择目标字段并按所有字段分组,然后计数。
当我检查 UI 时,发生了以下事情。
输入 81.2 GiB 随机写入 645.7 GiB
shuffle 写入的数据怎么会比原来读取的数据大这么多呢? 在这种情况下应该稍微扩展一下。 谁能解释一下?谢谢。
【问题讨论】:
-
这是因为 parquet 是经过编码和压缩的。数据膨胀后发生洗牌
-
我参考了这个答案stackoverflow.com/questions/48847660/… 并进行了一些优化,如下所示: select... .repartition(1000, $"field1") .sortWithinPartitions() .group... 新输入和随机写入数据为:input 40.2Gib,shuffle write 77.3Gib,shuffle write/input总是2左右。比未优化的40.7 vs. 334.9,比例为8要好很多。shuffle数据应该还是parquet+snappy,但是数据的组织方式可能会影响随机写入数据的大小。也许吧。
标签: apache-spark