Scala中ParallelTestExecution引发LiveListenerBus已停止的偶现测试失败
解决Spark-testing-base并行单测试套件时的LiveListenerBus已停止异常
问题场景
我在Scala项目中使用spark-testing-base做Spark测试,这个库明确建议不要设置parallelExecution in Test := true,因为Spark运行在单JVM环境下,多测试会共享资源引发冲突。但我尝试用Scala的ParallelTestExecution来并行单个测试套件内的测试时,出现了偶发的失败,报错如下:
java.lang.IllegalStateException: LiveListenerBus is stopped. [info] at org.apache.spark.scheduler.LiveListenerBus.addToQueue(LiveListenerBus.scala:92) [info] at org.apache.spark.scheduler.LiveListenerBus.addToStatusQueue(LiveListenerBus.scala:75) [info] at org.apache.spark.sql.internal.SharedState.<init>(SharedState.scala:115) [info] at org.apache.spark.sql.SparkSession.$anonfun$sharedState$1(SparkSession.scala:143) [info] at scala.Option.getOrElse(Option.scala:189) [info] at org.apache.spark.sql.SparkSession.sharedState$lzycompute(SparkSession.scala:143) [info] at org.apache.spark.sql.SparkSession.sharedState(SparkSession.scala:142) [info] at org.apache.spark.sql.SparkSession.$anonfun$sessionState$2(SparkSession.scala:162) [info] at scala.Option.getOrElse(Option.scala:189) [info] at org.apache.spark.sql.SparkSession.sessionState$lzycompute(SparkSession.scala:160)
我清楚这是因为测试线程尝试访问已停止的Spark上下文导致的。为了排查上下文停止的触发点,我在BeforeAll中添加了一个Spark监听器:
spark.sparkContext.addSparkListener(new SparkListener { override def onApplicationEnd(applicationEnd: SparkListenerApplicationEnd): Unit = { val stackTrace = Thread.currentThread().getStackTrace log.error(s"Spark application is shutting down for: ${stackTrace}") } })
添加这个监听器后,测试不再失败,而且每个测试结束后都会输出上面的日志,符合spark-testing-base中AfterAll的预期行为。
原异常原因分析
- spark-testing-base的设计逻辑是:每个测试套件共享一个SparkSession/SparkContext,在
BeforeAll阶段初始化,AfterAll阶段销毁。 - 当启用
ParallelTestExecution并行套件内的测试时,多个测试线程会同时操作这个共享的Spark上下文。 - 偶发失败的核心是线程时序冲突:某个测试线程执行完毕后,触发了
AfterAll的上下文销毁逻辑,但此时还有其他测试线程正在尝试初始化或使用SparkSession的SharedState——而SharedState初始化时会往LiveListenerBus中添加状态事件,一旦LiveListenerBus已经被停止,就会直接抛出IllegalStateException。 - 这种失败是偶发的,完全取决于JVM的线程调度时机,只有当销毁线程和使用线程的执行时序刚好重叠时才会触发。
添加监听器后问题解决的缘由
- 自定义
SparkListener的onApplicationEnd回调会在Spark上下文关闭时被执行,这个回调的执行会占用线程资源,相当于引入了一个微小的同步延迟,改变了线程的执行时序,避免了销毁线程刚关闭LiveListenerBus,就有其他线程立刻尝试往里面添加事件的极端冲突场景。 - 另外,Spark的
LiveListenerBus在存在自定义监听器时,关闭流程会等待所有监听器的回调执行完毕再完成最终的停止操作,这给剩余的测试线程足够的时间完成当前的资源使用,或者阻止了上下文被过早销毁。 - 本质上,这个监听器的添加消除了之前线程时序冲突的可能性,让测试的执行流程变得更稳定。
内容的提问来源于stack exchange,提问作者Adithya Aravind
相关产品推荐
相关产品推荐

