我想修复/合并我的数据,以便将其保存到每个分区的一个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有更好的方法吗?
我有一个相当简单的工作将日志文件转换为镶木地板.它正在处理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).
有时,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不降低并行度吗?
我正在尝试使用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数据的分区将作为单个文件被写出。
我最好的前进道路是什么?
当写入 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在前面添加
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。
我们可以写100个文件的数据,每个文件有10个分区吗?
我知道我们可以使用重新分区或合并来减少分区数量。但是我看到一些hadoop生成的avro数据具有比文件数量更多的分区。