在IntelliJ中运行Apache Spark单元测试出现SessionStateBuilder实例化错误
错误触发原因
这个错误的根因是SparkContext关联的LiveListenerBus已经被终止,导致新的SparkSession初始化SessionState失败,具体触发逻辑如下:
- 代码中存在两套SparkSession管理逻辑:一套是类全局初始化的
spark实例,通过getOrCreate()创建;另一套是DataTestUtils.withSpark工具方法内部创建的临时session。 withSpark这类工具方法的标准实现是:执行闭包逻辑后主动调用spark.stop()终止SparkContext,而getOrCreate()创建的SparkSession默认会复用同一个JVM内的SparkContext。- 当
withSpark执行完毕关闭Context后,后续再调用全局spark实例的方法(或者下一个测试用例复用这个全局实例),就会因为Context关联的LiveListenerBus已经停止触发报错。
解决方案
可任选以下一种方案修改:
方案1:统一使用withSpark管理SparkSession生命周期
把所有用到SparkSession的逻辑(包括读取csv的逻辑)全部移到withSpark的闭包内,用工具提供的session完成所有操作,不再单独维护全局spark实例:
"simple unit test" should "check for data correctness" in { appCfgT match { case Success(appCfg) => preStart() DataTestUtils.withSpark { session => // 把读csv的逻辑移到闭包里,用withSpark提供的session val rawDF: DataFrame = session .read .format("csv") .option("delimiter", ",") .option("timestampFormat", "yyyy/MM/dd HH:mm:ss ZZ") .option("inferSchema", value = true) .option("mode", "DROPMALFORMED") .option("header", value = true) .option("multiLine", value = true) .schema(encodedHousingSchema) .load(appCfg.sourceFileUrl) val rows = session.sparkContext.parallelize(Seq(new HousingModel())) val data = session.createDataFrame(rows) val verificationResult = VerificationSuite() .onData(data) .addCheck( Check(CheckLevel.Error, "unit testing my data") .hasSize(_ == 4092) .isComplete("id") .isUnique("id") .isComplete("productName") .isContainedIn("priority", Array("high", "low")) .isNonNegative("numViews") .containsURL("description", _ >= 0.5) .hasApproxQuantile("numViews", 0.5, _ <= 10) ) .run() } case Failure(fail) => fail("配置加载失败", fail) } }
方案2:统一使用全局SparkSession,调整withSpark实现
如果要保留全局的SparkSession,修改DataTestUtils.withSpark的逻辑:优先复用已存在的SparkSession,执行完闭包后不主动调用stop(),在测试类的afterAll钩子中统一关闭全局SparkSession。
全局初始化代码可补充配置避免多Context冲突:
val spark: SparkSession = SparkSession.builder() .config("spark.master", "local[*]") .config("spark.driver.allowMultipleContexts", "false") .appName("housing-data-test") .getOrCreate() // 测试类添加afterAll统一关闭 override def afterAll(): Unit = { spark.stop() super.afterAll() }
内容的提问来源于stack exchange,提问作者joesan
相关产品推荐
相关产品推荐

