You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark作业OOM且spark_driver_dagscheduler相关指标持续上升排查

指标含义解释

spark_driver_dagscheduler_messageprocessingtime_type_timers_count是Spark Driver端DAG调度器的消息处理时间计时器计数指标,持续暴涨说明:

  • DAG调度器接收到的待处理消息量在不断堆积,消息处理速度远低于消息生产速度
  • 大概率是Driver端出现了阻塞型任务或者内存泄漏,导致调度线程无法及时消费调度队列里的消息

内存泄漏定位及问题根因

从提供的代码可直接定位到以下核心问题:

  1. 重复创建SparkSession导致资源泄漏
    你在streamMqttMessageForSpecificColumn方法里每次调用都重新创建SparkSession,而不是复用main方法里已经初始化的实例。每启动一个流查询就多建一个SparkSession实例,所有Session对应的调度资源、状态对象都在Driver内存里堆积,很快就会占满内存。
  2. 静态CSV数据的持久化没有释放
    你用spark.read.csv读取的静态CSV数据调用了persist(StorageLevel.MEMORY_ONLY_SER),但静态数据和流数据关联之后,你从来没有主动调用unpersist释放缓存,每个查询都会持有这份缓存数据的引用,内存不会被自动回收。
  3. 流关联没有正确配置状态保留规则
    你用了流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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.29 07:54:04