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)
如何将数据集仅保存到一个实木复合地板文件中?谢谢
我的想法是使用 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,表示批处理间隔。
我希望您可以根据您的用例调整解决方案。
| 归档时间: |
|
| 查看次数: |
893 次 |
| 最近记录: |