相关疑难解决方法(0)

DataFrame partitionBy到单个Parquet文件(每个分区)

我想修复/合并我的数据,以便将其保存到每个分区的一个Parquet文件中.我还想使用Spark SQL partitionBy API.所以我可以这样做:

df.coalesce(1).write.partitionBy("entity", "year", "month", "day", "status")
  .mode(SaveMode.Append).parquet(s"$location")
Run Code Online (Sandbox Code Playgroud)

我已经测试了这个并且它似乎表现不佳.这是因为在数据集中只有一个分区可以处理,文件的所有分区,压缩和保存都必须由一个CPU内核完成.

在调用coalesce之前,我可以重写这个来手动执行分区(使用带有不同分区值的过滤器).

但是使用标准的Spark SQL API有更好的方法吗?

apache-spark apache-spark-sql

41
推荐指数
2
解决办法
4万
查看次数

Spark数据帧写入方法编写了许多小文件

我有一个相当简单的工作将日志文件转换为镶木地板.它正在处理1.1TB的数据(分为64MB - 128MB文件 - 我们的块大小为128MB),大约有12000个文件.

工作如下:

 val events = spark.sparkContext
  .textFile(s"$stream/$sourcetype")
  .map(_.split(" \\|\\| ").toList)
  .collect{case List(date, y, "Event") => MyEvent(date, y, "Event")}
  .toDF()

df.write.mode(SaveMode.Append).partitionBy("date").parquet(s"$path")
Run Code Online (Sandbox Code Playgroud)

它使用通用模式收集事件,转换为DataFrame,然后写出镶木地板.

我遇到的问题是,这会在HDFS集群上造成一些IO爆炸,因为它试图创建如此多的小文件.

理想情况下,我想在分区'date'中只创建一些镶木地板文件.

控制它的最佳方法是什么?是通过使用'coalesce()'吗?

这将如何影响给定分区中创建的文件数量?它取决于我在Spark中使用了多少执行程序?(目前设定为100).

scala apache-spark

9
推荐指数
2
解决办法
1万
查看次数

合并减少整个阶段的并行性(火花)

有时,Spark会以低效的方式“优化”数据框架计划。考虑以下Spark 2.1中的示例(也可以在Spark 1.6中复制):

val df = sparkContext.parallelize((1 to 500).map(i=> scala.util.Random.nextDouble),100).toDF("value")

val expensiveUDF = udf((d:Double) => {Thread.sleep(100);d})

val df_result = df
.withColumn("udfResult",expensiveUDF($"value"))

df_result
.coalesce(1)
.saveAsTable(tablename)
Run Code Online (Sandbox Code Playgroud)

在此示例中,我想在对数据帧进行昂贵的转换后写入1个文件(这只是一个演示此问题的示例)。Spark向上移动coalesce(1),使得UDF仅应用于包含1个分区的数据帧,从而破坏了并行性(有趣的repartition(1)是,这种行为不起作用)。

概括地说,当我想在转换的某个部分中增加并行度,但此后降低并行度时,就会发生此行为。

我发现了一种解决方法,包括缓存数据框,然后触发对数据框的完整评估:

val df = sparkContext.parallelize((1 to 500).map(i=> scala.util.Random.nextDouble),100).toDF("value")

val expensiveUDF = udf((d:Double) => {Thread.sleep(100);d})

val df_result = df
.withColumn("udfResult",expensiveUDF($"value"))
.cache

df_result.rdd.count // trigger computation

df_result
.coalesce(1)
.saveAsTable(tablename)
Run Code Online (Sandbox Code Playgroud)

我的问题是:在这种情况下,还有另一种方法可以告诉Spark不降低并行度吗?

scala apache-spark

6
推荐指数
1
解决办法
2066
查看次数

使用Spark的partitioningBy方法对S3中的大型偏斜数据集进行分区

我正在尝试使用Spark将较大的分区数据集写到磁盘上,并且该partitionBy算法在我尝试过的两种方法中都遇到了麻烦。

分区严重偏斜-有些分区很大,有些很小。

问题1:

当我之前使用repartition时repartitionBy,Spark将所有分区写为单个文件,即使是大文件也是如此

val df = spark.read.parquet("some_data_lake")
df
  .repartition('some_col).write.partitionBy("some_col")
  .parquet("partitioned_lake")
Run Code Online (Sandbox Code Playgroud)

这需要永远执行,因为Spark不会并行编写大型分区。如果其中一个分区具有1TB的数据,Spark将尝试将整个1TB的数据作为单个文件写入。

问题2:

当我不使用时repartition,Spark会写出太多文件。

此代码将写出疯狂的文件。

df.write.partitionBy("some_col").parquet("partitioned_lake")
Run Code Online (Sandbox Code Playgroud)

我在一个很小的8 GB数据子集上运行了此操作,Spark写入了85,000+个文件!

当我尝试在生产数据集上运行此文件时,一个包含1.3 GB数据的分区被写为3,100个文件。

我想要什么

我希望每个分区都写成1 GB文件。因此,具有7 GB数据的分区将作为7个文件被写出,而具有0.3 GB数据的分区将作为单个文件被写出。

我最好的前进道路是什么?

partitioning apache-spark apache-spark-sql

6
推荐指数
3
解决办法
631
查看次数

保存到分区 parquet 文件时实现并发

当写入 a dataframeto parquet使用时partitionBy:

df.write.partitionBy("col1","col2","col3").parquet(path)
Run Code Online (Sandbox Code Playgroud)

我期望正在写入的每个分区都是由单独的任务独立完成的,并且与分配给当前 Spark 作业的工作人员数量并行。

然而,在写入镶木地板时,实际上一次只有一个工作程序/任务在运行。该工作人员循环遍历每个分区并串行写出文件.parquet。为什么会出现这种情况 - 有没有办法强制执行此spark.write.parquet操作并发?

下面的不是我想看到的(应该是700%+..)

在此输入图像描述

从另一篇文章中我也尝试repartition在前面添加

Spark parquet 分区:大量文件

df.repartition("col1","col2","col3").write.partitionBy("col1","col2","col3").parquet(path)
Run Code Online (Sandbox Code Playgroud)

不幸的是,这没有任何效果:仍然只有一名工人..

local注意:我在模式下运行,local[8]并且看到其他Spark 操作与多达 8 个并发工作线程一起运行,并且使用了高达 750% 的 cpu。

scala apache-spark parquet

5
推荐指数
1
解决办法
5691
查看次数

使用少于N个分区的N个文件将数据写入磁盘

我们可以写100个文件的数据,每个文件有10个分区吗?

我知道我们可以使用重新分区或合并来减少分区数量。但是我看到一些hadoop生成的avro数据具有比文件数量更多的分区。

partition apache-spark

1
推荐指数
1
解决办法
3310
查看次数