Spark Structured Streaming内存溢出问题:定时存Parquet后如何清理数据?
兄弟,我太懂你这种扛着高流量流处理还遇到OOM的崩溃感了!咱们先拆解下问题,再一步步解决:
先搞清楚内存不足的核心原因
你说Spark DataFrame是“无界”的,但其实在Structured Streaming里,它是基于微批处理的——默认情况下,每个微批处理完数据后,只要没有状态留存,内存应该会自动释放。你遇到OOM,大概率是踩了这几个坑:
- 用了
Complete输出模式:这种模式会把全量聚合结果一直存在内存里,每批都更新全量,数据越攒越多肯定炸 - 状态操作没设过期:比如做窗口聚合、join这类带状态的操作,状态数据会一直存在内存里,没及时清理
- 内存配置没踩对点:executor的内存不光是堆内存,还有堆外内存的开销没考虑到
针对性解决方案,按优先级来
先检查输出模式,能换
Append就换
如果你的业务不需要全量结果,只是要把每5分钟的数据写入Parquet,那一定要用Append模式:streamingDF.writeStream .format("parquet") .option("path", "/your/output/path") .option("checkpointLocation", "/your/checkpoint/path") .outputMode("append") // 关键!只输出新增数据,不会留存全量 .trigger(Trigger.ProcessingTime("5 minutes")) // 刚好匹配你每5分钟存一次的需求 .start()Append模式下,每个微批处理完的数据写入后,对应的内存就会被释放,不会一直堆积。如果必须用状态操作,强制状态过期
要是你有窗口聚合、会话化这类需要保留状态的场景,一定要加withWatermark来设置状态过期时间,让Spark自动清理过期状态:val streamingDF = spark.readStream .format("your-source") .load() .withWatermark("event_time", "15 minutes") // 水位线设为比窗口大一点的时间 .groupBy(window($"event_time", "5 minutes"), $"key") .agg(count("*").as("cnt")) .writeStream .format("parquet") .outputMode("append") .trigger(Trigger.ProcessingTime("5 minutes")) .start()这里水位线设15分钟,意思是Spark会自动清理超过15分钟的事件对应的状态,不会让状态无限膨胀。
调整Executor内存配置,别只盯着堆内存
你给每个executor配了12G内存,但要注意:spark.executor.memoryOverhead:堆外内存,默认是executor内存的10%,但高流量场景下可能不够,建议设为4G(--conf spark.executor.memoryOverhead=4096),避免堆外内存溢出spark.memory.fraction:默认0.6,控制堆内存中用于存储(缓存、状态)和计算的比例,如果状态多,可以调到0.7,给存储更多空间(--conf spark.memory.fraction=0.7)
别手动缓存流DataFrame,必要时记得unpersist
有些同学会习惯性缓存流DataFrame,但Structured Streaming的微批机制本身会处理内存释放,手动缓存反而会导致数据堆积。如果真的有多次复用的场景,用完一定要调用.unpersist()。排查数据倾斜
100Mb/s的流如果某个Key特别热,会导致单个executor负载过高,内存爆掉。可以通过Spark UI的Stages页面看Task的执行时间和数据量,要是发现某个Task数据量是其他的几十倍,就得用加盐(给Key加随机前缀)的方式打散数据,处理完再聚合回来。
关于“保存后清理DataFrame”的误区
其实Structured Streaming里根本不需要你手动清理DataFrame——只要是Append模式,微批处理完成并写入后,该批次的DataFrame数据会被Spark的垃圾回收机制自动回收。你之前觉得“没法清理”,大概率是因为模式不对或者状态没过期,导致内存里一直留存着不必要的数据。
内容的提问来源于stack exchange,提问作者carl

