尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

SparkSession与SparkContext:从API演进到统一数据处理入口

SparkSession与SparkContext:从API演进到统一数据处理入口 1. 项目概述从SparkContext到SparkSession的演进之路如果你是从Spark 1.x时代一路用过来的老用户看到SparkSession这个API时可能会有点懵我们熟悉的那个SparkContext去哪了这两个家伙到底是什么关系为什么新写的代码里SparkSession几乎无处不在这其实是一个典型的“API演进”故事背后反映的是Spark从一个纯粹的分布式计算引擎向一个统一的数据处理与分析平台转变的宏大叙事。简单来说SparkSession是Spark 2.0引入的一个全新的、统一的入口点它整合了SparkContext、SQLContext、HiveContext以及StreamingContext等多个旧API的功能旨在提供一个更简洁、更一致的编程接口。理解它们的关系不仅是掌握Spark API变迁的关键更是写出更现代化、更健壮Spark代码的基础。无论你是刚入门的新手还是想升级旧代码的老手这篇文章都会带你彻底理清这两者的来龙去脉、核心差异以及最佳实践。2. 核心关系与设计理念拆解2.1 为什么需要SparkSession旧API的痛点在Spark 1.x时代开发者需要根据不同的数据处理需求与多个不同的“上下文”Context对象打交道。这带来了几个显著的痛点入口点分散如果你想做基础的RDD操作你需要创建SparkContext。如果你要使用DataFrame和SQL你需要一个SQLContext或者功能更强的HiveContext。如果你还要做流处理那你又需要一个StreamingContext。一个应用里管理多个上下文不仅代码冗余还增加了资源管理和生命周期的复杂性。配置不一致每个上下文对象都有自己的配置方式虽然底层共享同一个SparkConf但在API层面缺乏统一。比如在SQLContext中设置Hive支持和在SparkContext中设置序列化器感觉像是在操作两套不同的系统。会话Session概念的缺失对于交互式数据分析比如在Spark Shell或Notebook中用户可能需要隔离不同的计算会话每个会话有自己的临时表、UDF注册和配置。旧的API模型对此支持较弱。不利于统一的数据源访问随着Spark逐渐成为统一的数据处理平台需要一个更高层级的抽象来统一访问各种数据源如Hive表、Parquet文件、JSON、JDBC等并管理相关的元数据如临时视图。SparkSession的设计正是为了解决这些问题。它不是一个简单的包装器而是一个深思熟虑的、面向用户的新抽象层。2.2 SparkSession与SparkContext的包含关系最核心的关系可以概括为一个SparkSession实例内部封装了一个SparkContext实例。你可以把SparkSession想象成一个功能更全面的“总经理”而SparkContext是他手下专管核心计算资源的“技术总监”。在代码层面这种关系非常清晰。当你创建一个SparkSession后你可以通过其属性直接访问内部的SparkContext// 创建SparkSession val spark SparkSession.builder() .appName(MyApp) .master(local[*]) .getOrCreate() // 从SparkSession中获取其内部的SparkContext val sc: SparkContext spark.sparkContext // 现在你可以用sc来做传统的RDD操作 val rdd sc.textFile(data.txt)同样你也可以获取到SQLContext在SparkSession中就是它自己和StreamingContext需要通过SparkContext来创建。这意味着通过SparkSession这一个入口你获得了访问所有Spark核心功能的钥匙。注意在同一个JVM进程中SparkContext的设计是单例的。SparkSession虽然可以创建多个例如为不同的用户会话但在标准的Spark应用非交互式多会话场景中通常也遵循“一个应用一个SparkSession”的最佳实践。通过SparkSession.builder().getOrCreate()方法可以确保在同一个进程中获取到同一个实例避免资源冲突。2.3 新旧API功能映射与对比为了更直观地理解我们用一个表格来对比新旧API的核心功能功能模块Spark 1.x / 旧APISpark 2.x / SparkSession API说明应用入口与配置new SparkContext(conf)SparkSession.builder().config(conf).getOrCreate()SparkSession.Builder提供了链式调用的配置方式更优雅。RDD操作sc.textFile(),sc.parallelize()spark.sparkContext.textFile(),spark.sparkContext.parallelize()RDD API通过sparkContext属性访问完全保留。DataFrame SQLsqlContext.read.json(),sqlContext.sql()spark.read.json(),spark.sql()SparkSession自身就实现了SQLContext的接口直接调用。Hive支持new HiveContext(sc)SparkSession.builder().enableHiveSupport().getOrCreate()通过Builder的一个方法启用创建后即可操作Hive表。临时视图管理df.registerTempTable(tempView)(旧方法)df.createOrReplaceTempView(tempView)新方法名更准确视图生命周期与会话绑定。运行时配置sc.getConf或sqlContext.setConfspark.conf.set(spark.sql.shuffle.partitions, 200)提供了统一的conf属性来获取和设置Spark所有运行时配置。Catalyst优化器与Tungsten引擎需通过SQLContext间接使用内建于SparkSession对所有DataFrame/Dataset操作自动生效用户无需关心底层优化享受统一的高性能执行。从这个对比可以看出SparkSession将过去分散的能力集中到了一个简洁的API之下。对于开发者而言学习成本和代码维护成本都显著降低。3. 核心细节解析与实操要点3.1 SparkSession的创建与配置详解创建SparkSession的标准方式是使用Builder模式。这是你开始任何一个Spark 2.x应用的第一步也是最关键的一步。import org.apache.spark.sql.SparkSession val spark SparkSession.builder() .appName(My Spark Application) // 设置应用名称会在Web UI和日志中显示 .master(local[*]) // 设置运行模式local本地测试 yarnYARN集群 spark://host:portStandalone集群 .config(spark.sql.shuffle.partitions, 200) // 设置具体的Spark配置参数 .config(spark.executor.memory, 2g) .enableHiveSupport() // 如果需要操作Hive元数据库和表必须启用此项 .getOrCreate() // 核心方法获取已存在的Session或创建新的关键点解析getOrCreate()方法这是精髓所在。它会检查当前JVM进程中是否已经存在一个全局默认的SparkSession。如果存在则直接返回如果不存在则根据Builder的配置创建一个新的。这在交互式环境如Spark Shell、Zeppelin、Jupyter Notebook中尤其重要可以防止用户重复创建导致资源浪费和冲突。在独立的Spark应用中它通常就是创建新的。.master()的设定在本地测试时local[*]表示使用所有可用的CPU核心。在提交到集群时这个参数通常会被spark-submit命令行中的--master参数覆盖所以生产代码中有时会省略以增加灵活性。.enableHiveSupport()这个调用不是必须的但如果你需要使用Hive的元数据仓库即CREATE TABLE语句创建的表存储在Hive Metastore中。读写Hive格式的表。使用HiveQL特有的语法或UDF。 那么就必须启用它。启用后SparkSession会在内部创建一个HiveMetastoreClient来连接元数据库。注意启用Hive支持需要将Hive的相关Jar包放入classpath如果只是使用Spark内置的Derby内存数据库来创建临时元数据则不一定需要。实操心得配置的优先级Spark配置的加载是有优先级的理解这个可以避免配置不生效的坑。优先级从高到低通常是在代码中通过.config()直接设置的参数。提交应用时通过spark-submit的--conf参数传递的。应用JAR包中的spark-defaults.conf文件。Spark安装目录下的$SPARK_HOME/conf/spark-defaults.conf。系统环境变量。代码中的.config()具有最高优先级这让你可以在代码里写死一些关键配置如序列化方式确保其不会被外部配置意外覆盖。3.2 如何正确访问底层上下文对象虽然SparkSession是推荐的主要入口但某些场景下你依然需要直接操作底层的上下文对象。SparkSession提供了清晰的访问路径。// 1. 获取SparkContext - 用于RDD API、累加器、广播变量等 val sc spark.sparkContext // 示例创建RDD设置累加器 val accum sc.longAccumulator(MyAccumulator) val dataRdd sc.parallelize(1 to 100) dataRdd.foreach(x accum.add(x)) println(sAccumulator value: ${accum.value}) // 2. 获取SparkSession自身作为SQLContext - 实际上就是spark本身 // spark.read 等价于 spark.sqlContext.read // spark.sql(...) 等价于 spark.sqlContext.sql(...) val df spark.read.option(header, true).csv(people.csv) df.createOrReplaceTempView(people) val resultDF spark.sql(SELECT name, age FROM people WHERE age 20) // 3. 获取StreamingContext (对于Spark Streaming应用) // 注意Structured Streaming的入口是spark本身即spark.readStream // 这里指的是旧的DStream API import org.apache.spark.streaming._ val ssc new StreamingContext(sc, Seconds(1)) // 需要传入SparkContext和批处理间隔 // ... 后续DStream操作 // 4. 访问运行时配置 val shufflePartitions spark.conf.get(spark.sql.shuffle.partitions) println(sCurrent shuffle partitions: $shufflePartitions) // 动态修改配置某些配置在运行时可以修改 spark.conf.set(spark.sql.shuffle.partitions, 100)注意事项生命周期管理在Spark应用中SparkSession和SparkContext的生命周期应该与整个应用保持一致。通常在应用的main函数开始处创建在结束处调用spark.stop()来关闭。spark.stop()方法会优雅地停止内部的SparkContext释放所有资源如Executor进程、网络连接等。在长时间运行的服务如Spark Streaming应用中需要确保在收到关闭信号时正确调用stop()。3.3 临时视图Temporary View的会话隔离性这是SparkSession引入的一个非常重要的特性完美体现了“Session”会话的概念。// 假设我们有一个DataFrame val df spark.createDataFrame(Seq((Alice, 25), (Bob, 30))).toDF(name, age) // 创建一个临时视图它的生命周期与创建它的SparkSession绑定 df.createOrReplaceTempView(people_temp) // 可以在SQL中查询这个视图 spark.sql(SELECT * FROM people_temp).show() // 尝试在另一个新的SparkSession中查询这个视图会失败 val sparkNewSession SparkSession.builder() .appName(NewSession) .master(local[*]) .getOrCreate() // sparkNewSession.sql(SELECT * FROM people_temp) // 这会抛出异常Table or view not found // 创建全局临时视图Global Temporary View df.createOrReplaceGlobalTempView(people_global) // 全局临时视图被注册到全局数据库global_temp中跨Session可见 sparkNewSession.sql(SELECT * FROM global_temp.people_global).show() // 这样可以成功核心区别临时视图Temp View作用域仅限于创建它的那个SparkSession。当该Session停止后视图自动删除。非常适合在同一个应用或交互会话中的中间数据共享。全局临时视图Global Temp View作用域跨Session但绑定到同一个Spark应用即共享同一个SparkContext的所有SparkSession。它被保存在一个名为global_temp的全局数据库中查询时需要加上库名前缀。适用于需要在同一个应用内不同模块或线程间共享数据的场景。这个设计使得Spark可以更好地支持多用户交互环境如Thrift Server每个用户连接拥有自己独立的SparkSession他们的临时表互不干扰。4. 实操过程从旧代码迁移到新API4.1 迁移步骤与代码对比假设我们有一段经典的Spark 1.6代码它从文本文件创建RDD进行WordCount同时又将一个JSON文件读成DataFrame进行查询。我们来看如何将其升级到使用SparkSession。Spark 1.6 旧风格代码// 旧代码需要多个Context import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.sql.SQLContext object OldStyleApp { def main(args: Array[String]): Unit { // 1. 创建配置 val conf new SparkConf().setAppName(OldApp).setMaster(local[*]) // 2. 创建SparkContext (RDD入口) val sc new SparkContext(conf) // 3. 创建SQLContext (DataFrame入口) val sqlContext new SQLContext(sc) try { // 使用SparkContext做RDD操作 val textRDD sc.textFile(input.txt) val wordCounts textRDD .flatMap(_.split( )) .map(word (word, 1)) .reduceByKey(_ _) wordCounts.saveAsTextFile(output_rdd) // 使用SQLContext做DataFrame操作 val df sqlContext.read.json(people.json) df.createOrReplaceTempView(people) val adults sqlContext.sql(SELECT name FROM people WHERE age 18) adults.show() } finally { // 需要手动停止Context sc.stop() } } }Spark 2.x 新风格代码// 新代码统一使用SparkSession import org.apache.spark.sql.SparkSession object NewStyleApp { def main(args: Array[String]): Unit { // 1. 使用Builder模式创建SparkSession (统一入口) val spark SparkSession.builder() .appName(NewApp) .master(local[*]) .getOrCreate() // 2. 需要时从SparkSession中获取SparkContext val sc spark.sparkContext try { // RDD操作通过sc进行 (API保持不变) val textRDD sc.textFile(input.txt) val wordCounts textRDD .flatMap(_.split( )) .map(word (word, 1)) .reduceByKey(_ _) wordCounts.saveAsTextFile(output_rdd) // DataFrame操作直接通过spark进行 (更简洁) val df spark.read.json(people.json) // 注意是 spark.read不是 sqlContext.read df.createOrReplaceTempView(people) val adults spark.sql(SELECT name FROM people WHERE age 18) // 注意是 spark.sql adults.show() } finally { // 只需要停止SparkSession它会负责关闭内部的SparkContext spark.stop() } } }迁移要点总结导入包将import org.apache.spark.{SparkConf, SparkContext}和import org.apache.spark.sql.SQLContext替换为import org.apache.spark.sql.SparkSession。创建入口将分别创建SparkConf、SparkContext、SQLContext的步骤合并为使用SparkSession.builder()链式调用。对象引用将所有sqlContext.xxx如sqlContext.read,sqlContext.sql替换为spark.xxx。保留sc.xxxRDD操作但sc需要通过spark.sparkContext获取。资源关闭只需调用spark.stop()无需再单独调用sc.stop()。4.2 在结构化流Structured Streaming中的应用SparkSession在Structured Streaming中扮演了绝对核心的角色它统一了批处理和流处理的API。val spark SparkSession.builder() .appName(StructuredStreamingExample) .master(local[*]) .getOrCreate() // 读取流式数据源返回一个DataFrame即流式DataFrame val lines spark.readStream .format(socket) // 数据源格式socket, kafka, file等 .option(host, localhost) .option(port, 9999) .load() // 进行类似批处理的转换操作 val words lines.as[String].flatMap(_.split( )) val wordCounts words.groupBy(value).count() // 定义输出接收器并启动流查询 val query wordCounts.writeStream .outputMode(complete) // 输出模式complete, append, update .format(console) // 输出目的地console, memory, kafka等 .start() query.awaitTermination() // 等待流处理终止关键点你会发现除了readStream和writeStream中间的数据处理逻辑flatMap,groupBy与批处理完全一致。这正是SparkSession带来的统一编程模型的威力——批流一体。开发者无需学习两套不同的API。4.3 在多线程环境下的使用在Web服务或一些自定义的调度器中你可能会在多个线程中使用Spark。这里有一个重要的原则SparkContext和由同一个SparkContext创建的SparkSession不是线程安全的。// 错误示例在多个线程中共享同一个SparkSession对象进行并发操作 object UnsafeSparkApp extends App { val spark SparkSession.builder().appName(Unsafe).master(local[*]).getOrCreate() val df spark.range(10) (1 to 5).foreach { i new Thread(() { // 并发调用Spark操作可能导致不可预知的结果或错误 df.filter(sid $i).count() println(sThread $i done) }).start() } Thread.sleep(5000) spark.stop() }上述代码可能导致序列化错误、任务提交混乱等问题。正确的做法是每个线程使用独立的SparkSession但注意SparkSession.builder().getOrCreate()在同一个JVM内默认返回同一个实例。你需要使用newSession()方法来创建与父Session共享SparkContext但配置独立的新Session。object SafeSparkApp extends App { val parentSpark SparkSession.builder().appName(Safe).master(local[*]).getOrCreate() (1 to 5).foreach { i new Thread(() { // 为每个线程创建一个新的SparkSession它们共享底层的SparkContext val threadSpark parentSpark.newSession() try { val df threadSpark.range(10) df.filter(sid $i).count() println(sThread $i done) } finally { // 可以关闭线程Session但不会关闭底层的SparkContext // threadSpark.close() } }).start() } Thread.sleep(5000) parentSpark.stop() // 最终关闭父Session和SparkContext }newSession()创建的是一个轻量级的会话它会复制父Session的配置但拥有独立的SQL配置、临时视图、UDF注册等。底层SparkContext是共享的所以资源开销很小。将Spark操作封装为任务提交到Spark集群执行这是更常见的生产模式。主程序Driver创建SparkSession然后将数据处理逻辑封装到闭包中通过RDD.map、DataFrame.map等算子分发到Executor上执行。Executor上的代码是线程安全的由Spark框架管理。5. 常见问题与排查技巧实录在实际使用中从SparkContext过渡到SparkSession或者混合使用时会遇到一些典型问题。5.1 问题1java.lang.NoSuchMethodError或ClassNotFoundException问题描述在迁移旧项目或引入新依赖时运行时报错找不到SparkSession相关类或方法。原因分析这几乎总是依赖冲突或版本不匹配导致的。例如你的项目依赖的Spark Core版本是1.6但代码中却尝试导入Spark 2.0的SparkSession类。排查与解决检查构建工具依赖仔细检查你的pom.xmlMaven或build.sbtSBT文件确保所有Spark相关依赖spark-core,spark-sql,spark-hive等的版本号完全一致并且与你运行的Spark集群版本匹配。!-- Maven 示例确保版本号统一 -- dependency groupIdorg.apache.spark/groupId artifactIdspark-core_2.12/artifactId version3.3.0/version !-- 核心版本 -- scopeprovided/scope !-- 通常集群已提供打包时排除 -- /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-sql_2.12/artifactId version3.3.0/version !-- 必须与core版本一致 -- scopeprovided/scope /dependency使用依赖树分析工具使用mvn dependency:tree或sbt dependencyTree命令查看完整的依赖关系检查是否有其他第三方库传递依赖了不同版本的Spark Jar包。如果有需要使用exclusion规则排除。确认运行环境如果你是用spark-submit提交任务确保提交命令中--packages参数指定的版本或者环境变量SPARK_HOME指向的Spark安装目录版本与代码编译版本一致。5.2 问题2启用Hive支持失败问题描述调用.enableHiveSupport()后应用启动失败抛出类似org.apache.spark.sql.AnalysisException: org.apache.hadoop.hive.ql.metadata.HiveException: java.lang.RuntimeException: Unable to instantiate org.apache.hadoop.hive.ql.metadata.SessionHiveMetaStoreClient;的异常。原因分析Spark需要连接一个Hive Metastore服务来持久化表元数据。当你启用Hive支持时Spark默认会尝试连接一个运行中的Metastore如通过hive.metastore.uris配置。如果没配置或连接失败并且你没有在代码中创建任何需要持久化的表只使用临时视图可能不会报错。但如果你尝试执行CREATE TABLE等操作或者Spark内部需要访问Metastore时就会失败。解决方案仅使用内存元数据开发/测试如果你不需要真正的Hive Metastore只是想使用HiveQL语法和UDF可以配置Spark使用Derby内存数据库。确保spark.sql.warehouse.dir指向一个本地可写目录Spark 2.0并且没有配置hive.metastore.uris。这样Spark会在本地启动一个内嵌的Derby实例。val spark SparkSession.builder() .appName(TestHive) .master(local[*]) .config(spark.sql.warehouse.dir, /tmp/spark-warehouse) // 本地目录 // 不设置 hive.metastore.uris .enableHiveSupport() .getOrCreate()连接远程Hive Metastore生产在生产环境中你需要正确配置Metastore连接。val spark SparkSession.builder() .appName(ProdHive) .config(hive.metastore.uris, thrift://metastore-host:9083) .enableHiveSupport() .getOrCreate()同时确保提交任务的机器上有正确的Hive-site.xml配置文件在classpath中通常放在$SPARK_HOME/conf/目录下该文件包含了Metastore连接等详细信息。检查依赖确保你的应用包含了Hive相关的依赖如spark-hive_2.12并且版本匹配。5.3 问题3临时视图“找不到表”问题描述在同一个应用里在一个地方创建了临时视图df.createOrReplaceTempView(my_table)但在另一个地方用spark.sql(SELECT * FROM my_table)查询时却报错Table or view my_table not found。原因分析这通常是因为你使用了不同的SparkSession实例。临时视图是绑定到创建它的特定SparkSession实例的。如果你无意中创建了另一个SparkSession例如在某个函数内部又调用了一次SparkSession.builder().getOrCreate()并且因为上下文不同而创建了新实例那么在新Session中自然看不到旧Session创建的视图。排查与解决传递Session引用最佳实践是在应用的入口处创建一次SparkSession然后将其作为参数显式传递给所有需要用到它的函数或类。避免在各个模块内部自行创建。object MainApp { def processData(spark: SparkSession): Unit { val df spark.read.json(...) df.createOrReplaceTempView(temp1) // ... 其他操作 } def main(args: Array[String]): Unit { val spark SparkSession.builder()...getOrCreate() processData(spark) // 将spark实例传进去 spark.stop() } }使用全局临时视图如果确实需要在不同上下文中共享一个视图并且这些上下文属于同一个Spark应用共享SparkContext可以考虑使用createOrReplaceGlobalTempView。查询时需要使用global_temp.view_name。检查getOrCreate的调用上下文确保在不需要新Session的地方getOrCreate()调用能正确地获取到全局实例。在简单的单线程应用中这通常不是问题。5.4 问题4性能调优配置的生效位置问题描述我知道一些性能调优参数比如spark.sql.shuffle.partitions但应该在哪里设置在SparkSession.builder().config()里设置还是在创建后通过spark.conf.set()设置有什么区别原因与解决方案这两种方式在大多数情况下效果是等价的但有一个细微的时机差别。在Builder中设置.config()这些配置在SparkSession及其内部的SparkContext初始化时就生效了。这对于那些必须在SparkContext启动前就确定的参数是唯一的设置方式例如spark.serializer,spark.executor.memory,spark.master等。在创建后设置spark.conf.set()这些配置在Session创建后动态修改。对于Spark SQL相关的配置如spark.sql.shuffle.partitions,spark.sql.autoBroadcastJoinThreshold这种方式是有效的因为SQL优化器在规划每个查询时会读取当前的配置值。但是对于Spark Core的某些配置在Context启动后修改可能不会生效。最佳实践将所有的静态配置放在Builder里尤其是集群资源相关executor内存/核心数、序列化、shuffle服务等核心配置。这保证了应用启动时环境就是正确的。val spark SparkSession.builder() .appName(MyApp) .master(yarn) .config(spark.executor.memory, 4g) .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .config(spark.sql.shuffle.partitions, 200) // SQL配置放这里也没问题 .getOrCreate()将需要根据数据动态调整的SQL配置放在创建后设置例如在读取数据后你根据数据量大小动态调整广播join的阈值。val df spark.read.parquet(huge_table.parquet) val estimatedSize ... // 估算df大小 if (estimatedSize threshold) { spark.conf.set(spark.sql.autoBroadcastJoinThreshold, estimatedSize 1) }使用spark-submit的--conf参数对于生产部署将配置外部化是更好的选择这样无需修改代码即可调整参数。代码中的配置可以作为默认值或强制值。理解SparkSession和SparkContext的关系远不止是记住几个API调用那么简单。它意味着你接受了Spark向更高层次抽象和统一API演进的设计哲学。在绝大多数新项目中你应该毫不犹豫地将SparkSession作为唯一的起点。对于遗留代码逐步迁移到新API不仅能提升代码的简洁性和一致性也能让你更自然地运用结构化APIDataFrame/Dataset和结构化流Structured Streaming这些更现代、性能更好的组件。当你下次再看到SparkSession时希望你能清晰地认识到它不仅仅是SparkContext的替代品更是通往Spark统一数据处理世界的大门。
返回列表