Spark结构化流写入实木复合地板会创建许多文件

tan*_*oup 5 apache-spark parquet spark-streaming

我使用结构化流从kafka加载消息,进行一些聚合,然后将其写入镶木地板文件。问题是,仅从kafka发送的100条消息就创建了太多的实木复合地板文件(800个文件)。

聚合部分是:

return model
            .withColumn("timeStamp", col("timeStamp").cast("timestamp"))
            .withWatermark("timeStamp", "30 seconds")
            .groupBy(window(col("timeStamp"), "5 minutes"))
            .agg(
                count("*").alias("total"));
Run Code Online (Sandbox Code Playgroud)

查询:

StreamingQuery query = result //.orderBy("window")
            .writeStream()
            .outputMode(OutputMode.Append())
            .format("parquet")
            .option("checkpointLocation", "c:\\bigdata\\checkpoints")
            .start("c:\\bigdata\\parquet");
Run Code Online (Sandbox Code Playgroud)

当使用spark加载一个实木复合地板文件时,它显示为空

+------+-----+
|window|total|
+------+-----+
+------+-----+
Run Code Online (Sandbox Code Playgroud)

如何将数据集仅保存到一个实木复合地板文件中?谢谢

Dav*_*ein 2

我的想法是使用 Spark 结构化流处理来自 Azure Even Hub 的事件,然后以镶木地板格式将它们存储在存储上。

我终于弄清楚如何处理创建的许多小文件。Spark 版本 2.4.0。

这就是我的查询的样子

dfInput
  .repartition(1, col('column_name'))
  .select("*")
  .writeStream
  .format("parquet")
  .option("path", "adl://storage_name.azuredatalakestore.net/streaming")
  .option("checkpointLocation", "adl://storage_name.azuredatalakestore.net/streaming_checkpoint")
  .trigger(processingTime='480 seconds')
  .start()
Run Code Online (Sandbox Code Playgroud)

因此,我每 480 秒就会在存储位置创建一个文件。要找到文件大小和文件数量之间的平衡以避免 OOM 错误,只需使用两个参数:分区数量和processingTime,表示批处理间隔。

我希望您可以根据您的用例调整解决方案。