【发布时间】:2020-08-24 17:17:01
【问题描述】:
问题大纲:假设我在 AWS 的 EMR 集群上使用 spark 处理了 300+ GB 的数据。该数据具有三个属性,用于在文件系统上进行分区以在 Hive 中使用:日期、小时和(比方说)anotherAttr。我想以尽量减少写入文件数量的方式将此数据写入 fs。
我现在正在做的是获取日期、小时、anotherAttr 的不同组合以及组成组合的行数。我将它们收集到驱动程序上的列表中,并遍历列表,为每个组合构建一个新的 DataFrame,使用行数重新分区该 DataFrame 以推测文件大小,并使用 DataFrameWriter 将文件写入磁盘,.orcfinish关掉。
出于组织原因,我们不使用 Parquet。
这种方法效果不错,解决了下游团队使用 Hive 而不是 Spark 看不到由于大量文件导致的性能问题的问题。例如,如果我使用整个 300 GB 数据帧,对 1000 个分区(在 spark 中)和相关列进行重新分区,然后将其转储到磁盘,它会并行转储,并在大约 9 分钟内完成整个事情。但这会为较大的分区提供多达 1000 个文件,这会破坏 Hive 的性能。或者它会破坏某种性能,老实说不是 100% 确定是什么。我刚刚被要求尽量减少文件数量。使用我使用的方法,我可以将文件保持在我想要的任何大小(无论如何都相对接近),但没有并行性,运行大约需要 45 分钟,主要是等待文件写入。
在我看来,由于某些源行和某些目标行之间存在一对一的关系,并且由于我可以将数据组织到不重叠的“文件夹”(Hive 的分区)中,所以我应该能够以这样一种方式组织我的代码/数据帧,我可以要求 spark 并行编写所有目标文件。有人对如何攻击这个有建议吗?
我测试过的东西不起作用:
使用 Scala 并行集合启动写入。无论 spark 对 DataFrame 做什么,它都没有很好地分离任务,并且一些机器遇到了大量的垃圾收集问题。
DataFrame.map - 我试图映射唯一组合的 DataFrame,并从那里开始写入,但无法从
map中访问我实际需要的数据的 DataFrame -执行器上的 DataFrame 引用为 null。DataFrame.mapPartitions - 一个非初学者,无法从 mapPartitions 中提出任何想法来做我想做的事情
“分区”这个词在这里也不是特别有用,因为它既指火花按某些标准分割数据的概念,也指数据在磁盘上为 Hive 组织的方式。我想我在上面的用法中很清楚。因此,如果我想为这个问题提供一个完美的解决方案,那就是我可以创建一个基于三个属性的具有 1000 个分区的 DataFrame 以进行快速查询,然后从中创建另一个 DataFrame 集合,每个 DataFrame 都具有一个独特的组合这些属性,重新分区(在 spark 中,但对于 Hive),分区数与其包含的数据大小相适应。大多数 DataFrame 将有 1 个分区,少数将有多达 10 个。文件应该是 ~3 GB,我们的 EMR 集群的 RAM 比每个执行程序的 RAM 多,所以我们不应该看到这些文件对性能造成影响“大”分区。
一旦创建了 DataFrame 列表并且每个都重新分区,我可以要求 spark 将它们全部并行写入磁盘。
spark 中是否可能发生这样的事情?
我在概念上不清楚的一件事:说我有
val x = spark.sql("select * from source")
和
val y = x.where(s"date=$date and hour=$hour and anotherAttr=$anotherAttr")
和
val z = x.where(s"date=$date and hour=$hour and anotherAttr=$anotherAttr2")
y 在多大程度上与z 是不同的DataFrame?如果我重新分区y,那么洗牌对z 和x 有什么影响?
【问题讨论】:
-
佩服你的坚持,看看能不能学。我也喜欢赏金!
标签: parallel-processing apache-spark-sql orc