如何在Spark Shell中切换活跃的Spark Session?
如何切换Spark Shell中的活跃Spark Session
嗨,我来帮你搞定这个问题!你观察得很对——用spark.newSession()创建的新Session不会自动成为当前活跃的那个,默认操作还是会绑定到初始的spark实例上。不过有几种靠谱的方法可以切换活跃Session,让后续操作都基于新Session来执行:
1. 显式调用新Session的API
最直接的方式就是别用默认的toDF()这类无上下文的方法,而是直接通过新Session来创建DataFrame、注册视图。这样创建出来的对象自然就和新Session绑定了:
val spark2 = spark.newSession() // 用spark2来创建DataFrame,而不是直接用Seq.toDF() val s2 = spark2.createDataFrame(Seq(4,5,6)).toDF("num") // 验证一下,这个DataFrame属于spark2 s2.sparkSession // 会返回你的spark2实例(比如org.apache.spark.sql.SparkSession@43e869ea)
2. 在Spark Shell里重新赋值默认的spark变量
因为Spark Shell是交互式环境,你可以直接把新Session赋值给默认的spark变量,这样后续所有默认操作都会自动用这个新Session:
val spark2 = spark.newSession() // 把新Session赋值给默认的spark变量 val spark = spark2 // 现在再创建DataFrame,就会关联到spark2了 val s = Seq(1,2,3).toDF("num") s.sparkSession // 返回的就是spark2的实例
小提示:这种方式只适合Spark Shell这类交互式场景,在提交的批量Spark应用里别这么干,容易造成代码上下文混乱。
3. 使用官方的setActiveSession方法(Spark 2.4及以上版本)
从Spark 2.4开始,官方提供了SparkSession.setActiveSession()方法,这是最规范的切换活跃Session的方式,不管是Shell还是应用程序都能用:
val spark2 = spark.newSession() // 把spark2设置为当前活跃Session SparkSession.setActiveSession(spark2) // 现在创建DataFrame、注册临时视图都会默认用spark2了 val s = Seq(1,2,3).toDF("num") s.sparkSession // 指向的就是spark2
最后再补充一点:每个SparkSession都有自己独立的临时视图、UDF、配置上下文,切换活跃Session后,之前Session里的临时视图在新Session里是看不到的,这点和你之前观察到的行为一致,属于正常现象。
内容的提问来源于stack exchange,提问作者Somesh Dhal
相关产品推荐
相关产品推荐

