Spark作业OOM且spark_driver_dagscheduler相关指标持续上升排查
指标含义解释
spark_driver_dagscheduler_messageprocessingtime_type_timers_count是Spark Driver端DAG调度器的消息处理时间计时器计数指标,持续暴涨说明:
- DAG调度器接收到的待处理消息量在不断堆积,消息处理速度远低于消息生产速度
- 大概率是Driver端出现了阻塞型任务或者内存泄漏,导致调度线程无法及时消费调度队列里的消息
内存泄漏定位及问题根因
从提供的代码可直接定位到以下核心问题:
- 重复创建SparkSession导致资源泄漏
你在streamMqttMessageForSpecificColumn方法里每次调用都重新创建SparkSession,而不是复用main方法里已经初始化的实例。每启动一个流查询就多建一个SparkSession实例,所有Session对应的调度资源、状态对象都在Driver内存里堆积,很快就会占满内存。 - 静态CSV数据的持久化没有释放
你用spark.read.csv读取的静态CSV数据调用了persist(StorageLevel.MEMORY_ONLY_SER),但静态数据和流数据关联之后,你从来没有主动调用unpersist释放缓存,每个查询都会持有这份缓存数据的引用,内存不会被自动回收。 - 流关联没有正确配置状态保留规则
你用了流rate表和静态CSV表做关联,虽然加了watermark,但关联条件里没有针对流状态的清理规则,每个批次的状态会一直累加在Driver端。
修复建议
- 去掉
streamMqttMessageForSpecificColumn方法里的SparkSession创建逻辑,直接传入main方法中已经初始化的SparkSession实例 - 静态CSV数据读取后不需要持久化,或者在所有流查询启动完成后主动调用
unpersist释放缓存 - 因为是流和静态表的关联,可以改为将CSV数据广播后做广播关联,避免状态累加:
rate.join(broadcast(cvsStream), expr("csv.id == mod(counter.value,10)"))
内容的提问来源于stack exchange,提问作者Eljah
相关产品推荐
相关产品推荐

