【问题标题】:Splitting spark data into partitions and writing those partitions to disk in parallel将 spark 数据拆分为分区并将这些分区并行写入磁盘
【发布时间】: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 并行编写所有目标文件。有人对如何攻击这个有建议吗?

我测试过的东西不起作用:

  1. 使用 Scala 并行集合启动写入。无论 spark 对 DataFrame 做什么,它都没有很好地分离任务,并且一些机器遇到了大量的垃圾收集问题。

  2. DataFrame.map - 我试图映射唯一组合的 DataFrame,并从那里开始写入,但无法从 map 中访问我实际需要的数据的 DataFrame -执行器上的 DataFrame 引用为 null。

  3. 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,那么洗牌对zx 有什么影响?

【问题讨论】:

  • 佩服你的坚持,看看能不能学。我也喜欢赏金!

标签: parallel-processing apache-spark-sql orc


【解决方案1】:

此声明:

我将它们收集到驱动程序上的列表中,并遍历列表, 为每个组合构建一个新的 DataFrame,重新分区 DataFrame 使用行数来推测文件大小,以及 使用 DataFrameWriter 将文件写入磁盘,.orc 完成。

在涉及 Spark 的情况下完全偏离光束。收集到驱动程序从来都不是一个好方法,体积和 OOM 问题以及您的方法中的延迟很高。

使用以下方法可以简化并获得 Spark 的并行性,从而为您的老板节省时间和金钱:

df.repartition(cols...)...write.partitionBy(cols...)...

通过repartition 进行洗牌,partitionBy 不会洗牌。

就这么简单,使用 Spark 的默认并行性。

【讨论】:

  • 不幸的是,这会为 spark 在.repartition 步骤中创建的每个分区写入一个文件到磁盘。我试图在问题中概述这正是我要解决的问题。
  • 我试图向您概述您应该按预期使用软件。收集是不行的......
  • 我要求的是按预期使用软件的其他方式。发现一件你不喜欢我正在做的事情,然后建议我不要那样做,而是去做我寻求帮助而不做的事情,这没有帮助。
  • 没有,而且收集到驱动程序的方法很奇怪,因为它违反了所有 Spark 原则 afaik。没有帮助但很现实。这就是生活。我是建筑师,但它不会通过。
【解决方案2】:

我们(几乎)遇到了同样的问题,最终我们直接使用 RDD(而不是 DataFrames)并实现了我们自己的分区机制(通过扩展 org.apache.spark.Partitioner)

详细信息:我们正在读取来自 Kafka 的 JSON 消息。 JSON 应按 customerid/date/more 字段分组,并使用 Parquet 格式在 Hadoop 中编写,不要创建太多小文件。

步骤是(简化版): a) 从 Kafka 读取消息并将其转换为 RDD[(GroupBy, Message)] 的结构。 GroupBy 是一个案例类,包含用于分组的所有字段。

b) 使用 reduceByKeyLocally 转换并获取每个组的指标映射(消息数量/消息大小/等) - 例如 Map[GroupBy, GroupByMetrics]

c) 创建一个 GroupPartitioner,它使用之前收集的指标(以及一些输入参数,如所需的 Parquet 大小等)来计算应该为每个 GroupBy 对象创建多少个分区。基本上我们正在扩展 org.apache.spark.Partitioner 并覆盖 numPartitions 和 getPartition(key: Any)

d)我们使用之前定义的分区器对 a) 中的 RDD 进行分区:newPartitionedRdd = rdd.partitionBy(ourCustomGroupByPartitioner)

e) 使用两个参数调用 spark.sparkContext.runJob:第一个是在 d) 分区的 RDD,第二个是自定义函数 (func: (TaskContext, Iterator[T]) 将写入获取的消息从 Iterator[T] 到 Hadoop/Parquet

假设我们有 1 亿条消息,像这样分组

Group1 - 2 百万

Group2 - 8000 万

Group3 - 1800 万 我们决定每个分区必须使用 150 万条消息才能获得大于 500MB 的 Parquet 文件。我们最终会得到 Group1 的 2 个分区,Group2 的 54 个分区,Group3 的 12 个分区。

【讨论】:

  • 问题: 1. 您是否通过多次执行步骤 e 来实现并行性?在您的示例中,您是否会执行三次,每组一次,每次都使用具有不同分区计数的自定义分区器的不同实例化?换句话说,您在步骤 d 是否有一个串行瓶颈,类似于我如何“获取日期、小时、anotherAttr 和组成组合的行数的不同组合”?
  • 2.您编写消息的客户职能部门是否一次编写一条消息?我在 s3 中给 orc 写信,我一直在使用 DataFrameWriter。如果我必须使用 Iterator[T],我可能必须先在 EMR 中“本地”写入磁盘,然后再将文件复制到 s3。如果您对该自定义功能有任何更具体的建议,我会全力以赴!
  • 1 - 不,我只执行一次 runJob,在点 d) 创建/分区的 RDD 上。唯一可能的瓶颈可能是分区完全不平衡,但实际上我们没有这种情况。这个周末我会在我的 github 位置发布一个示例。
  • 2 - 在我们的例子中,我们逐行处理并转换为 Parquet,将结果写入内存并在 Azure ADLS 中关闭((或内存已满)时刷新。我们在这方面投入了一些时间,我们不得不对各种内部 Spark 类进行影子处理(我想同样的方法也可以用于 ORC 生成)。性能不错(优于默认的 Spark Dataframe->Parquet 生成),但我们为此投入了一些时间。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-05
  • 1970-01-01
  • 1970-01-01
  • 2016-12-13
  • 2018-05-24
相关资源
最近更新 更多