如何在Spark中分区和写入DataFrame而不删除没有新数据的分区?

jay*_*son 25 partitioning apache-spark parquet spark-dataframe

我试图DataFrame用Parquet格式保存到HDFS,使用DataFrameWriter三列值进行分区,如下所示:

dataFrame.write.mode(SaveMode.Overwrite).partitionBy("eventdate", "hour", "processtime").parquet(path)
Run Code Online (Sandbox Code Playgroud)

正如提到的这个问题,partitionBy将在删除分区的全部现有层次path,并在分区取而代之dataFrame.由于特定日期的新增量数据将定期出现,我想要的是仅替换层次结构中dataFrame具有数据的那些分区,而保持其他分区不变.

要做到这一点,似乎我需要使用其完整路径单独保存每个分区,如下所示:

singlePartition.write.mode(SaveMode.Overwrite).parquet(path + "/eventdate=2017-01-01/hour=0/processtime=1234567890")
Run Code Online (Sandbox Code Playgroud)

但是我无法理解将数据组织到单分区中DataFrame的最佳方法,以便我可以使用它们的完整路径将它们写出来.一个想法是这样的:

dataFrame.repartition("eventdate", "hour", "processtime").foreachPartition ...
Run Code Online (Sandbox Code Playgroud)

但foreachPartition操作上Iterator[Row]是不理想的写出来镶木格式.

我还考虑使用a select...distinct eventdate, hour, processtime获取分区列表,然后按每个分区过滤原始数据帧并将结果保存到完整的分区路径.但是,每个分区的独特查询加过滤器似乎效率不高,因为它会进行大量的过滤/写入操作.

我希望有一种更简洁的方法来保留dataFrame没有数据的现有分区?

谢谢阅读.

Spark版本:2.1

rod*_*mbs 17

这是一个老话题,但我遇到了同样的问题并找到了另一个解决方案,只需使用以下命令将分区覆盖模式设置为动态:

spark.conf.set('spark.sql.sources.partitionOverwriteMode', 'dynamic')
Run Code Online (Sandbox Code Playgroud)

因此,我的 spark 会话配置如下:

spark = SparkSession.builder.appName('AppName').getOrCreate()
spark.conf.set('spark.sql.sources.partitionOverwriteMode', 'dynamic')
Run Code Online (Sandbox Code Playgroud)

  • 这应该被标记为真正的解决方案。也许它比较慢,但它满足了OP的要求。 (2认同)
  • 在 Databricks 9.1 LTS(包括 Apache Spark 3.1.2、Scala 2.12)上有效并且没有看到性能下降 (2认同)

new*_*dev 12

模式选项Append有一个问题!

df.write.partitionBy("y","m","d")
.mode(SaveMode.Append)
.parquet("/data/hive/warehouse/mydbname.db/" + tableName)
Run Code Online (Sandbox Code Playgroud)

我已经测试并发现这将保留现有的分区文件.但是,这次出现的问题如下:如果您运行相同的代码两次(使用相同的数据),那么它将创建新的镶木地板文件,而不是替换相同数据的现有代码(Spark 1.6).因此,Append我们仍然可以解决这个问题,而不是使用Overwrite.我们应该在分区级别覆盖,而不是在表级覆盖.

df.write.mode(SaveMode.Overwrite)
.parquet("/data/hive/warehouse/mydbname.db/" + tableName + "/y=" + year + "/m=" + month + "/d=" + day)
Run Code Online (Sandbox Code Playgroud)

有关更多信息,请参阅以下链接:

在spark数据帧写入方法中覆盖特定分区

(我在suriyanto的评论之后更新了我的回复.Tnnx.)

  • 您是否测试过两次写入相同数据时它是否替换旧分区?从我的测试中,它实际上在分区目录中创建了一个新的镶木地板文件,导致数据加倍.我在Spark 2.2上. (3认同)

D3V*_*D3V 5

我知道这已经很老了。由于我看不到任何已发布的解决方案,因此我将继续发布一个。这种方法假设您在要写入的目录上有一个 hive 表。处理此问题的一种方法是创建一个临时视图dataFrame,应将其添加到表中,然后使用正常的类似 hive 的insert overwrite table ...命令:

dataFrame.createOrReplaceTempView("temp_view")
spark.sql("insert overwrite table table_name partition ('eventdate', 'hour', 'processtime')select * from temp_view")
Run Code Online (Sandbox Code Playgroud)

它保留旧分区,同时(覆盖)仅写入新分区。