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

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版本的一个已知问题,根源在于缓存管理模块在清理缓存时的并发安全漏洞:

  1. 当流任务长时间运行时,即使堆内存使用率不高,BlockManager管理的缓存块可能因为堆外内存不足或者LRU缓存驱逐策略被自动清理;
  2. 在unpersist()执行过程中,Spark的CacheManager会检查缓存块的加载状态,但此时部分块已经被回收,导致CachedRDDBuilder.isCachedRDDLoaded方法中访问了空对象,触发NullPointerException;
  3. 你的代码中在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 20:27:41