在sbt控制台运行Spark时任务执行停滞的问题排查
Spark在sbt控制台执行停滞问题的解决
问题现象
Spark新手,build.sbt配置正确(IntelliJ中代码可正常运行),sbt控制台能正常导入所需包,但直接执行以下代码时任务停滞无法完成:
import org.apache.spark._ val sc = new SparkContext("local[1]", "SimpleProg") val nums = sc.parallelize(List(1, 2, 3, 4)); println(nums.reduce((a, b) => a - b))
但将代码封装成函数后可正常运行,在spark-shell中执行也无问题。通过-Xprint:typer对比了三种场景的编译输出。
问题原因
sbt REPL在处理顶级作用域的RDD变量时,会自动调用ScalaRunTime.replStringOf生成变量的字符串展示,这个过程会触发RDD的惰性求值逻辑。在local[1]单线程模式下,REPL主线程等待RDD字符串化结果,而Spark的任务执行线程又依赖主线程资源,形成死锁,导致任务停滞。
从编译输出可验证:
- 停滞的sbt控制台输出中,
$eval的$print方法明确调用了replStringOf($line2.$read.INSTANCE.$iw.$iw.nums, 1000),尝试序列化RDD用于展示。 - spark-shell的REPL针对Spark对象做了特殊优化,嵌套的可序列化类结构避免了不必要的求值触发。
- 封装为函数后,
$result类型为Unit,$print返回空字符串,不会触发RDD的序列化操作,因此无阻塞。
解决方法
- 封装代码到函数/代码块:将Spark逻辑放入函数内执行,避免顶级变量触发REPL的自动展示逻辑:
def runSparkJob(): Unit = { val sc = new SparkContext("local[1]", "SimpleProg") val nums = sc.parallelize(List(1, 2, 3, 4)) println(nums.reduce((a, b) => a - b)) sc.stop() // 务必关闭资源 } runSparkJob() - 避免顶级作用域定义Spark对象:在sbt REPL中,不要直接在顶级定义
SparkContext、RDD等对象,尽量放在局部作用域内。 - 使用
sbt run执行代码:将测试代码写入Scala源文件,通过sbt run命令执行,绕开REPL的特殊处理逻辑。
附:各场景-Xprint:typer输出
1. 任务停滞的sbt控制台输出
[[syntax trees at end of typer]] // <console> package $line2 { object $read extends scala.AnyRef { def <init>(): type = { $read.super.<init>(); () }; object $iw extends scala.AnyRef { def <init>(): type = { $iw.super.<init>(); () }; object $iw extends scala.AnyRef { def <init>(): type = { $iw.super.<init>(); () }; import org.apache.spark._; private[this] val sc: org.apache.spark.SparkContext = new org.apache.spark.SparkContext("local[1]", "SimpleProg", spark.this.SparkContext.<init>$default$3, spark.this.SparkContext.<init>$default$4, spark.this.SparkContext.<init>$default$5); <stable> <accessor> def sc: org.apache.spark.SparkContext = $iw.this.sc; private[this] val nums: org.apache.spark.rdd.RDD[Int] = $iw.this.sc.parallelize[Int](scala.collection.immutable.List.apply[Int](1, 2, 3, 4), $iw.this.sc.parallelize$default$2[Nothing])((ClassTag.Int: scala.reflect.ClassTag[Int])); <stable> <accessor> def nums: org.apache.spark.rdd.RDD[Int] = $iw.this.nums } }; private[this] val INSTANCE: $line2.$read.type = this; <stable> <accessor> def INSTANCE: $line2.$read.type = $read.this.INSTANCE } } [[syntax trees at end of typer]] // <console> package $line2 { object $eval extends scala.AnyRef { def <init>(): $line2.$eval.type = { $eval.super.<init>(); () }; <stable> <accessor> lazy val $result: org.apache.spark.rdd.RDD[Int] = $line2.$read.INSTANCE.$iw.$iw.nums; <stable> <accessor> lazy val $print: String = { $line2.$read.INSTANCE.$iw.$iw; val sb: StringBuilder = new scala.`package`.StringBuilder(); sb.append("import org.apache.spark._\n"); sb.append("[1m[34msc[0m: [1m[32morg.apache.spark.SparkContext[0m = ".+(scala.runtime.ScalaRunTime.replStringOf($line2.$read.INSTANCE.$iw.$iw.sc, 1000))); sb.append("[1m[34mnums[0m: [1m[32morg.apache.spark.rdd.RDD[Int][0m = ".+(scala.runtime.ScalaRunTime.replStringOf($line2.$read.INSTANCE.$iw.$iw.nums, 1000))); sb.toString() } } }
2. 代码正常运行的spark-shell输出
[[syntax trees at end of typer]] // <console> package $line15 { sealed class $read extends AnyRef with java.io.Serializable { def <init>(): $line15.$read = { $read.super.<init>(); () }; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; private[this] val $line3$read: $line3.$read.type = $line3.$read.INSTANCE; <stable> <accessor> def $line3$read: $line3.$read.type = $iw.this.$line3$read; import $iw.this.$line3$read.$iw.$iw.spark; import $iw.this.$line3$read.$iw.$iw.sc; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; import org.apache.spark.SparkContext._; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; import $iw.this.$line3$read.$iw.$iw.spark.implicits._; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; import $iw.this.$line3$read.$iw.$iw.spark.sql; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; import org.apache.spark.sql.functions._; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; sealed class $iw extends AnyRef with java.io.Serializable { def <init>(): $iw = { $iw.super.<init>(); () }; private[this] val nums: org.apache.spark.rdd.RDD[Int] = $iw.this.$line3$read.$iw.$iw.sc.parallelize[Int](scala.collection.immutable.List.apply[Int](1, 2, 3, 4), $iw.this.$line3$read.$iw.$iw.sc.parallelize$default$2[Nothing])((ClassTag.Int: scala.reflect.ClassTag[Int])); <stable> <accessor> def nums: org.apache.spark.rdd.RDD[Int] = $iw.this.nums }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $iw.this.$iw(); <stable> <accessor> def $iw: $iw = $iw.this.$iw }; private[this] val $iw: $iw = new $read.this.$iw(); <stable> <accessor> def $iw: $iw = $read.this.$iw }; object $read extends scala.AnyRef with Serializable { def <init>(): type = { $read.super.<init>(); () }; private[this] val INSTANCE: $line15.$read = new $read(); <stable> <accessor> def INSTANCE: $line15.$read = $read.this.INSTANCE; <synthetic> private def readResolve(): Object = $line15.$read } }
3. 在sbt控制台封装为temp函数后正常运行的输出
[[syntax trees at end of typer]] // <console> package $line7 { object $read extends scala.AnyRef { def <init>(): type = { $read.super.<init>(); () }; object $iw extends scala.AnyRef { def <init>(): type = { $iw.super.<init>(); () }; object $iw extends scala.AnyRef { def <init>(): type = { $iw.super.<init>(); () }; import org.apache.spark._; object $iw extends scala.AnyRef { def <init>(): type = { $iw.super.<init>(); () }; import $line6.$read.INSTANCE.$iw.$iw.$iw.$iw.temp; object $iw extends scala.AnyRef { def <init>(): type = { $iw.super.<init>(); () }; private[this] val res1: Unit = $line6.$read.INSTANCE.$iw.$iw.$iw.$iw.temp(); <stable> <accessor> def res1: Unit = $iw.this.res1 } } } }; private[this] val INSTANCE: $line7.$read.type = this; <stable> <accessor> def INSTANCE: $line7.$read.type = $read.this.INSTANCE } } [[syntax trees at end of typer]] // <console> package $line7 { object $eval extends scala.AnyRef { def <init>(): $line7.$eval.type = { $eval.super.<init>(); () }; <stable> <accessor> lazy val $result: Unit = $line7.$read.INSTANCE.$iw.$iw.$iw.$iw.res1; <stable> <accessor> lazy val $print: String = { $line7.$read.INSTANCE.$iw.$iw.$iw.$iw; "" } } }
内容的提问来源于stack exchange,提问作者Dhruv
相关产品推荐
相关产品推荐

