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

在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("sc: org.apache.spark.SparkContext = ".+(scala.runtime.ScalaRunTime.replStringOf($line2.$read.INSTANCE.$iw.$iw.sc, 1000)));
      sb.append("nums: org.apache.spark.rdd.RDD[Int] = ".+(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 22:18:09