Spark 3.2.1流处理任务运行约1天后抛出NullPointerException问题排查咨询
问题分析:Spark 3.2.1流处理任务运行1天后抛出NullPointerException
问题背景
将Spark版本升级至3.2.1后,流处理任务运行约1天会抛出NullPointerException。当前集群配置为1个Driver和2个Executor,每个Executor分配2GB内存,老年代内存使用率约50%,堆内存状态看似健康,具体堆内存信息如下:
Heap par new generation total 307840K, used 239453K [0x0000000080000000, 0x0000000094e00000, 0x00000000aaaa0000) eden space 273664K, 81% used [0x0000000080000000, 0x000000008da4bdd0, 0x0000000090b40000) from space 34176K, 46% used [0x0000000092ca0000, 0x0000000093c2b6b8, 0x0000000094e00000) to space 34176K, 0% used [0x0000000090b40000, 0x0000000090b40000, 0x0000000092ca0000) concurrent mark-sweep generation total 811300K, used 451940K [0x00000000aaaa0000, 0x00000000dc2e9000, 0x0000000100000000) Metaspace used 102593K, capacity 110232K, committed 121000K, reserved 1155072K class space used 12473K, capacity 13482K, committed 15584K, reserved 1048576K
核心代码片段
此类流查询共有4个,核心逻辑如下(Scala代码):
sparkSession .readStream .format("kafka") .load .repartition(4) // ... 省略project和watermark逻辑 ... .watermark .groupby(k1, k2) .agg(size(collect_set("xxx"))) .writeStream .foreachBatch(test) .start val test: (Dataset[Row], Long) => Unit = (ds: Dataset[Row], _: Long) => { ds.persist(StorageLevel.MEMORY_AND_DISK_SER) ds.write .option("collection", "col_1") .option("maxBatchSize", "2048") .mode("append") .mongo() ds.write .option("collection", "col_2") .option("maxBatchSize", "2048") .mode("append") .mongo() ds.unpersist() }
异常栈信息
22/05/09 21:11:28 ERROR streaming.MicroBatchExecution: Query rydts_regist_gp [id = 669c2031-71b2-422b-859d-336722d289e9, runId = 049de32c-e6ff-48f1-8742-bb95122a36ea] terminated with error java.lang.NullPointerException at org.apache.spark.sql.execution.columnar.CachedRDDBuilder.$anonfun$isCachedRDDLoaded$1(InMemoryRelation.scala:248) at org.apache.spark.sql.execution.columnar.CachedRDDBuilder.$anonfun$isCachedRDDLoaded$1$adapted(InMemoryRelation.scala:247) at scala.collection.IndexedSeqOptimized.prefixLengthImpl(IndexedSeqOptimized.scala:41) at scala.collection.IndexedSeqOptimized.forall(IndexedSeqOptimized.scala:46) at scala.collection.IndexedSeqOptimized.forall$(IndexedSeqOptimized.scala:46) at scala.collection.mutable.ArrayOps$ofRef.forall(ArrayOps.scala:198) at org.apache.spark.sql.execution.columnar.CachedRDDBuilder.isCachedRDDLoaded(InMemoryRelation.scala:247) at org.apache.spark.sql.execution.columnar.CachedRDDBuilder.isCachedColumnBuffersLoaded(InMemoryRelation.scala:241) at org.apache.spark.sql.execution.CacheManager.$anonfun$uncacheQuery$8(CacheManager.scala:189) at org.apache.spark.sql.execution.CacheManager.$anonfun$uncacheQuery$8$adapted(CacheManager.scala:176) at scala.collection.TraversableLike.$anonfun$filterImpl$1(TraversableLike.scala:304) at scala.collection.Iterator.foreach(Iterator.scala:943) at scala.collection.Iterator.foreach$(Iterator.scala:943) at scala.collection.AbstractIterator.foreach(Iterator.scala:1431) at scala.collection.IterableLike.foreach(IterableLike.scala:74) at scala.collection.IterableLike.foreach$(IterableLike.scala:73) at scala.collection.AbstractIterable.foreach(Iterable.scala:56) at scala.collection.TraversableLike.filterImpl(TraversableLike.scala:303) at scala.collection.TraversableLike.filterImpl$(TraversableLike.scala:297) at scala.collection.AbstractTraversable.filterImpl(Traversable.scala:108) at scala.collection.TraversableLike.filter(TraversableLike.scala:395) at scala.collection.TraversableLike.filter$(TraversableLike.scala:395) at scala.collection.AbstractTraversable.filter(Traversable.scala:108) at org.apache.spark.sql.execution.CacheManager.recacheByCondition(CacheManager.scala:219) at org.apache.spark.sql.execution.CacheManager.uncacheQuery(CacheManager.scala:176) at org.apache.spark.sql.Dataset.unpersist(Dataset.scala:3220) at org.apache.spark.sql.Dataset.unpersist(Dataset.scala:3231) at common.job.xxx$.$anonfun$main$3(xxx.scala:117) at common.job.xxx$.$anonfun$main$3$adapted(xxx.scala:103) at org.apache.spark.sql.execution.streaming.sources.ForeachBatchSink.addBatch(ForeachBatchSink.scala:35) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$17(MicroBatchExecution.scala:600) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:103) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:163) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:90) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$16(MicroBatchExecution.scala:598) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:69) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runBatch(MicroBatchExecution.scala:598) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:228) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:69) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:193) at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:57) at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:187) at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:303) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775) at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:286) at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:209)
原因分析
这个NPE是Spark 3.2.x版本的一个已知问题,根源在于缓存管理模块在清理缓存时的并发安全漏洞:
- 当流任务长时间运行时,即使堆内存使用率不高,BlockManager管理的缓存块可能因为堆外内存不足或者LRU缓存驱逐策略被自动清理;
- 在
unpersist()执行过程中,Spark的CacheManager会检查缓存块的加载状态,但此时部分块已经被回收,导致CachedRDDBuilder.isCachedRDDLoaded方法中访问了空对象,触发NullPointerException; - 你的代码中在foreachBatch里对同一个Dataset进行两次写操作后调用
unpersist(),这种场景下缓存块的生命周期管理更容易触发这个bug。
解决方案
针对这个问题,你可以按照以下优先级尝试解决:
1. 升级Spark版本到3.3.0及以上
Spark官方在3.3.0版本中修复了这个缓存清理时的NPE问题,这是最彻底的解决方案。升级后,缓存模块的并发处理逻辑被优化,不会再出现此类空指针异常。
2. 调整unpersist()的调用方式
在你的代码中,将ds.unpersist()改为ds.unpersist(true),启用阻塞式清理:
ds.unpersist(true) // 等待缓存块被完全清理后再继续执行
这种方式可以避免并发清理时的空指针问题,确保缓存块的状态检查是在块完全存在或完全清理后进行。
3. 优化缓存策略与资源配置
- 调整StorageLevel:如果你的任务对磁盘IO不敏感,可以考虑改用
StorageLevel.DISK_ONLY_SER,减少内存缓存的压力,降低块被LRU驱逐的概率; - 增加堆外内存:虽然堆内存健康,但BlockManager默认会使用堆外内存存储缓存块。可以通过配置
spark.memory.offHeap.size增加堆外内存,比如设置为1g:spark.memory.offHeap.size=1g spark.memory.offHeap.enabled=true - 减小micro batch数据量:通过调整
spark.sql.streaming.microBatchInterval或者Kafka源的maxOffsetsPerTrigger参数,减少每个批次处理的数据量,降低缓存的负载。
4. 避免重复复用同一个Dataset缓存
如果业务允许,可以考虑将Dataset复制一份,分别用于两个MongoDB的写操作,避免对同一个缓存块进行多次访问和清理:
val ds1 = ds.persist(StorageLevel.MEMORY_AND_DISK_SER) ds1.write... // 写入col_1 ds1.unpersist(true) val ds2 = ds.persist(StorageLevel.MEMORY_AND_DISK_SER) ds2.write... // 写入col_2 ds2.unpersist(true)
这种方式虽然会增加一点计算开销,但可以避免缓存块的并发访问冲突。
内容的提问来源于stack exchange,提问作者cxb
相关产品推荐
相关产品推荐

