1. 从Spark 2.0的“统一入口”说起
如果你是从Spark 1.x时代一路走过来的老用户,第一次在Spark 2.0的代码里看到SparkSession时,心里多半会嘀咕一句:“这又是个啥?我的SparkContext和SQLContext呢?” 这种感觉就像你熟悉了家里的老式收音机,突然给你换了个智能音箱,虽然功能更强了,但一开始总有点不习惯。SparkSession的出现,正是为了解决这种“不习惯”,它本质上是一个为了简化开发体验而设计的统一入口。
在Spark 1.x时代,要开发一个应用,你常常需要和好几个“上下文”对象打交道。想用RDD?你得先创建一个SparkContext。想用DataFrame或SQL?那你还需要一个SQLContext(或者针对Hive的HiveContext)。这些对象各自为政,创建和管理起来既繁琐又容易出错,尤其是在需要共享配置或临时表的时候。SparkSession的诞生,就是为了终结这种“多入口”的混乱局面。它把SparkContext、SQLContext以及后续版本中引入的StreamingContext(在Structured Streaming中)的核心功能,都整合到了一个统一的API之下。
所以,当你现在写spark = SparkSession.builder.appName(“MyApp”).getOrCreate()时,你得到的不仅仅是一个会话,而是一个功能齐全的“瑞士军刀”。你可以通过spark.sparkContext来访问底层的RDD操作,通过spark.sql来执行SQL查询,通过spark.read和spark.write来轻松读写各种数据源。这种设计极大地简化了代码,也让Spark的API对新手更加友好。但这也带来了新的疑问:SparkSession和SparkContext现在到底是什么关系?是替代,是封装,还是共存?接下来,我们就一层层剥开来看。
2. SparkContext:分布式计算的“发动机”
要理解SparkSession,我们必须先回到基石——SparkContext。你可以把它想象成整个Spark应用程序的发动机和总指挥。它是Driver程序与集群资源管理器(如YARN、Mesos或Standalone Master)通信的桥梁,也是创建所有分布式数据集(RDD)和广播变量、累加器等共享变量的唯一入口。
2.1 SparkContext的核心职责
它的工作非常繁重,主要包括以下几个方面:
连接集群与资源申请:
SparkContext在初始化时,会根据配置(SparkConf)向集群资源管理器申请Executor资源。它负责与Cluster Manager谈判:“我需要多少个Executor,每个需要多少内存和CPU核心。” 一旦申请成功,它就负责与这些Executor保持心跳通信,指挥它们干活。创建RDD的工厂:所有RDD的创建,无论是通过
parallelize从本地集合生成,还是通过textFile从HDFS读取,亦或是通过hadoopFile访问Hadoop数据源,最终都需要SparkContext来执行。它定义了数据的分区逻辑和计算位置。DAG调度与任务分发:当你对RDD进行一系列转换操作(如
map、filter)时,SparkContext中的DAGScheduler会将这些操作翻译成一个有向无环图(DAG)。然后,TaskScheduler负责将这个DAG拆分成一个个具体的Task,并将这些Task分发到各个Executor上去执行。这是Spark高效运行的核心。管理共享变量:为了在分布式任务间高效共享数据,
SparkContext提供了**广播变量(Broadcast Variables)和累加器(Accumulators)**的创建接口。广播变量用于只读数据的缓存分发,累加器用于安全地聚合各任务的计算结果。作业与状态监控:通过
SparkContext,你可以获取当前应用ID(applicationId)、Web UI的URL、以及作业的执行状态。它是你洞察应用运行情况的窗口。
在代码层面,一个典型的Spark 1.x应用是这样开始的:
// Spark 1.x 风格 val conf = new SparkConf().setAppName(“MyApp”).setMaster(“local[*]”) val sc = new SparkContext(conf) val data = sc.textFile(“hdfs://path/to/data”) val words = data.flatMap(_.split(“ “))这里,sc就是一切的核心。没有它,你的Spark应用根本无法启动。
2.2 SparkContext的“单例”限制与多上下文困境
然而,SparkContext有一个非常重要的设计限制:在一个JVM进程中,只能有一个活跃的SparkContext实例。如果你尝试创建第二个,它会直接抛出异常。这个设计保证了集群资源管理的唯一性和一致性。
这个限制本身是合理的,但在Spark生态逐渐丰富后,带来了开发上的不便。随着DataFrame和SQL API的引入,SQLContext(及其子类HiveContext)成为了新的必需品。虽然SQLContext内部依赖于一个SparkContext,但它们毕竟是不同的对象。在同一个应用中,你可能需要同时操作RDD和DataFrame,代码里就需要同时维护sc和sqlContext两个引用。更麻烦的是,如果你想在不同的线程或模块中共享配置、临时表或者UDF(用户自定义函数),你需要小心翼翼地传递这些上下文对象,或者使用全局变量,这增加了代码的复杂度和耦合性。
3. SparkSession:新时代的“统一指挥中心”
正是为了解决上述痛点,Spark 2.0引入了SparkSession。它不是SparkContext的简单替代品,而是一个更高层次的抽象和封装,其核心目标是提供一套统一的、用户友好的API来访问Spark的所有功能。
3.1 SparkSession的“三位一体”架构
你可以把SparkSession看作一个“外壳”或“门面”(Facade Pattern),它内部封装并统一管理了多个重要的上下文对象。通过SparkSession的实例(通常命名为spark),你可以无缝访问到:
- SparkContext: 通过
spark.sparkContext属性访问。所有原有的RDD API依然在这里。 - SQLContext: 通过
spark.sqlContext属性访问。但更常见的是直接使用spark本身的方法,因为SparkSession直接集成了SQLContext的几乎所有功能。 - HiveContext(如果启用Hive支持): 当创建
SparkSession时通过.enableHiveSupport()方法启用后,spark就具备了完整的Hive元数据访问和HQL支持能力。 - StreamingContext(在Structured Streaming中): 对于Structured Streaming API,
SparkSession同样是入口,用于定义流式DataFrame。
这种设计带来了巨大的便利。看看Spark 2.0+的典型代码:
// Spark 2.0+ 风格 import org.apache.spark.sql.SparkSession val spark = SparkSession.builder .appName(“UnifiedExample”) .config(“spark.some.config”, “some-value”) .getOrCreate() import spark.implicits._ // 使用DataFrame API (原来需要SQLContext) val df = spark.read.json(“examples/src/main/resources/people.json”) df.show() // 使用SQL (原来需要SQLContext) df.createOrReplaceTempView(“people”) val sqlDF = spark.sql(“SELECT * FROM people”) sqlDF.show() // 使用RDD API (通过sparkContext,原来需要SparkContext) val rdd = spark.sparkContext.textFile(“examples/src/main/resources/people.txt”) rdd.collect().foreach(println)所有操作,通过一个spark对象全部搞定。代码更简洁,依赖更清晰。
3.2 SparkSession独有的高级功能
除了整合旧API,SparkSession还引入或强化了一些独有的特性,使其不仅仅是简单的包装:
统一的配置管理:通过
SparkSession.builder可以集中设置所有Spark配置,这些配置会对SparkContext、SQLContext等所有底层上下文生效,保证了配置的一致性。全局临时视图(Global Temporary View):这是
SparkSession一个非常重要的增强。在SQLContext中创建的临时视图(createTempView)是会话(Session)级别的,只在创建它的DataFrame所在的SQLContext中可见。而SparkSession允许创建全局临时视图(createGlobalTempView),这些视图被绑定到一个全局的global_temp数据库,可以在同一个Spark应用程序内的不同SparkSession实例间共享。这对于模块化应用或测试场景非常有用。更便捷的UDF注册:注册UDF(用户自定义函数)可以直接在
spark.udf命名空间下进行,更加直观。内置的
spark对象(在Spark Shell中):如果你使用spark-shell或pyspark交互式环境,你会发现一个预创建好的名为spark的SparkSession对象已经在那里等着你了,开箱即用,体验无缝。
4. 关系深度剖析:封装、依赖与生命周期
理解了各自角色后,我们来精确地定义它们的关系。
4.1 不是替代,而是演进与封装
最关键的结论是:SparkSession并没有取代SparkContext,而是将其作为核心组件封装在内。SparkContext依然是Spark运行时引擎的绝对核心,负责最底层的集群通信、任务调度和RDD管理。SparkSession是在此基础上构建的一个更友好、功能更全面的客户端API。
从生命周期上看,这种封装关系体现得非常明显:
- 当你调用
SparkSession.builder().getOrCreate()时,如果当前JVM内没有活跃的SparkContext,它会首先创建一个SparkContext(以及SQLContext)。 - 因此,一个
SparkSession实例必然对应一个内部的SparkContext实例。你可以通过spark.sparkContext获取到它。 - 反之则不成立。在Spark 2.0之前,你可以只有
SparkContext而没有SparkSession。在Spark 2.0+中,虽然你可以通过new SparkContext()的方式“单独”创建它(通常不推荐),但SparkSession的builder在检测到已有SparkContext存在时,会复用这个上下文,而不是创建新的。
4.2 依赖关系的代码级验证
我们可以写一段简单的代码来验证这种关系:
val spark = SparkSession.builder.appName(“Test”).master(“local”).getOrCreate() println(s“SparkSession is created: $spark”) println(s“Internal SparkContext: ${spark.sparkContext}”) println(s“Are they the same object? ${spark eq spark.sparkContext}”) // false, 它们是不同的对象 println(s“SparkContext‘s appName: ${spark.sparkContext.appName}”) // 应该输出 ‘Test‘ spark.stop() // 停止SparkSession // 此时再尝试访问 spark.sparkContext 会抛出异常,因为底层的SparkContext也被停止了。这段代码清晰地表明,spark和spark.sparkContext是两个不同的对象引用,但后者是前者内部状态的一部分。停止SparkSession会连带停止其内部的SparkContext。
4.3 何时该用哪个?
对于开发者来说,一个很实际的问题是:我该用哪个?
对于所有Spark 2.0+的新项目,无脑使用
SparkSession作为唯一入口。这是官方推荐的最佳实践。通过它,你可以访问Spark的所有功能(RDD, DataFrame, SQL, Streaming)。代码更干净,功能更全面。只有在极少数需要直接操作非常底层API的情况下,才需要通过
spark.sparkContext去访问SparkContext的原生方法。例如:- 使用一些尚未集成到DataFrame API中的特殊数据源。
- 操作累加器和广播变量(虽然
SparkSession也提供了快捷方式,但底层仍是SparkContext)。 - 获取一些底层运行时信息,如
applicationId、uiWebUrl等(同样,SparkSession也提供了sparkContext属性来访问)。
对于维护遗留的Spark 1.x代码,在升级到Spark 2.x时,一个常见的迁移路径就是将
SparkContext和SQLContext的创建逻辑,替换为创建一个SparkSession,然后通过这个session来获取原有的上下文对象,这样可以最小化代码改动。
5. 实战配置、调优与常见“坑点”
了解了理论,我们来看看在实际开发和运维中,围绕这两个对象有哪些需要注意的实操细节。
5.1 正确创建与配置SparkSession
创建SparkSession的最佳实践是使用Builder模式。getOrCreate()方法尤为重要,它保证了在同一个JVM内(例如,在某个Web服务中多次调用初始化代码)只会创建一个SparkSession实例,避免了资源浪费和冲突。
import org.apache.spark.sql.SparkSession val spark = SparkSession.builder .appName(“MyProductionJob”) // 设置应用名,在集群UI中显示 .master(“yarn”) // 或 “local[*]”, “spark://master:7077” // 动态设置配置,优先级高于spark-defaults.conf .config(“spark.sql.shuffle.partitions”, “200”) // 调整Shuffle分区数,对性能影响巨大 .config(“spark.executor.memory”, “4g”) .config(“spark.driver.memory”, “2g”) .config(“spark.serializer”, “org.apache.spark.serializer.KryoSerializer”) // 使用Kryo序列化提升性能 // 启用Hive支持,即可访问Hive元数据仓库和HQL语法 .enableHiveSupport() // 获取或创建实例 .getOrCreate() // 设置日志级别,避免输出过多INFO日志 spark.sparkContext.setLogLevel(“WARN”)注意:
.config()设置的参数会覆盖spark-defaults.conf配置文件中的默认值,但会被Spark提交脚本中的--conf参数或代码中通过SparkConf直接设置的值所覆盖。配置的优先级需要心中有数。
5.2 性能调优关联点
SparkSession和SparkContext的配置直接影响性能:
Executor资源:通过
.config(“spark.executor.memory”, …)和.config(“spark.executor.cores”, …)设置的资源,最终是由底层的SparkContext去和YARN等资源管理器申请的。申请不足会导致任务运行慢,申请过多则会造成集群资源浪费。Shuffle分区数:
spark.sql.shuffle.partitions(默认200)这个配置至关重要。它决定了Spark SQL或DataFrame操作进行Shuffle(如join, groupBy)时产生的分区数量。如果这个值设置得过大,会产生大量小任务,增加调度开销;设置得过小,则每个分区数据量过大,可能导致Executor内存溢出(OOM)。通常需要根据数据量大小进行调整,一般建议为executor-cores * executor-num的2-3倍。序列化:
.config(“spark.serializer”, “org.apache.spark.serializer.KryoSerializer”)。Kryo序列化比默认的Java序列化更快、序列化后的数据更小,能显著减少网络传输和内存占用。但需要注意注册自定义类(.registerKryoClasses),否则Kryo效率会下降。
5.3 开发者常踩的“坑”与解决方案
坑点一:在Spark Streaming (DStreams) 中误用
这是一个经典误区。传统的Spark Streaming(基于DStream的API)的入口是StreamingContext,它需要接收一个SparkContext作为参数。如果你已经有了一个SparkSession,正确的做法是复用其内部的SparkContext,而不是再创建一个。
// 正确做法 val spark = SparkSession.builder…getOrCreate() val ssc = new StreamingContext(spark.sparkContext, Seconds(1)) // 复用SparkContext // 错误做法 val spark = SparkSession.builder…getOrCreate() val sc = new SparkContext(…) // 尝试创建第二个SparkContext,会抛出异常! val ssc = new StreamingContext(sc, Seconds(1))坑点二:临时视图的作用域混淆
如前所述,createTempView创建的视图是会话级别的。如果你在某个函数内创建了一个SparkSession的副本,或者从已有的spark对象newSession()了一个子会话,那么在这个新会话中,你将看不到之前会话创建的普通临时视图。
val spark1 = SparkSession.builder…getOrCreate() df.createOrReplaceTempView(“table1”) val spark2 = spark1.newSession() // 创建一个新的会话 spark2.sql(“SELECT * FROM table1”) // 这里会报错:Table or view not found // 如果需要跨会话共享,请使用全局临时视图 df.createOrReplaceGlobalTempView(“global_table1”) spark2.sql(“SELECT * FROM global_temp.global_table1”) // 可以成功查询坑点三:在UDF或闭包中不当引用SparkSession/SparkContext
在RDD的map、filter等操作中,或者注册的UDF函数内部,如果引用了外部的SparkSession或SparkContext对象,这些对象需要被序列化并发送到Executor端。这常常会导致Task not serializable错误。
val spark = … // Driver端的SparkSession val broadcastVar = spark.sparkContext.broadcast(一些数据) // 正确:使用广播变量 val rdd = spark.sparkContext.parallelize(1 to 10) // 错误做法:在闭包中直接引用spark rdd.map { x => // 这里试图使用spark,但spark无法被序列化到Executor // spark.sql(...) // 这行会报序列化错误 x + broadcastVar.value // 正确:使用广播变量的值 } // 正确的UDF定义也应避免内部引用SparkSession spark.udf.register(“myUdf”, (x: Int) => x * 2) // 纯函数,安全坑点四:未正确关闭资源
在长时间运行的服务(如Thrift JDBC/ODBC Server)或单元测试中,如果反复创建SparkSession而不关闭,会导致资源(如端口绑定、内存)泄漏。务必使用try-finally或SparkSession的stop()方法确保资源释放。
val spark = SparkSession.builder…getOrCreate() try { // 你的业务逻辑 } finally { spark.stop() }6. 从源码角度理解设计
对于想深入理解的同学,可以简要看一下源码中的关系。在Spark源码的SparkSession类中,你可以找到如下关键字段:
// 摘自 Apache Spark 源码 (简化) class SparkSession private( @transient val sparkContext: SparkContext, @transient private val existingSharedState: Option[SharedState], …) extends Serializable with Closeable with Logging { // … private[sql] val sessionState: SessionState = … private[sql] lazy val sharedState: SharedState = … // … }可以看到,SparkSession的主构造函数中,sparkContext是一个必需的参数。在SparkSession的伴生对象builder的getOrCreate()方法中,逻辑是:先尝试获取或创建SparkContext,然后再用这个SparkContext作为参数来实例化SparkSession。这从源码层面证实了SparkSession对SparkContext的依赖关系。
SharedState和SessionState是另外两个关键内部类。SharedState持有跨会话共享的状态,如全局临时视图、共享的HiveClient等;而SessionState则持有会话级别的状态,如临时视图、UDF注册信息、SQL配置等。这解释了为什么不同SparkSession可以共享全局视图,而普通临时视图不行。
理解到这个层次,你就能真正明白,SparkSession是一个精心设计的、管理着多种会话和共享状态的高层管理器,而SparkContext是其麾下负责具体分布式计算执行的引擎主管。两者各司其职,共同构成了现代Spark应用开发的基石。