ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Zeppelin Spark Interpreter 模块架构与多版本支持深度解析

Apache Zeppelin Spark Interpreter 模块架构与多版本支持深度解析 数据分析数据可视化大数据后端前端任务调度【免费下载链接】zeppelinWeb-based notebook that enables>项目地址https://gitcode.com/gh_mirrors/zeppe/zeppelin点击查看免费下载Spark interpreter 是 Apache Zeppelin 中最早、也最核心的解释器它让用户可以围绕 Spark 会话直接在 Notebook 中运行 Scala、SQL 与 PySpark 代码。本文将围绕 spark/README.md 描述的核心模块结构展开结合spark模块下的真实源码与构建配置深入讲解 Spark interpreter 是如何通过Java 入口 多 Scala 版本适配 多 Spark 版本 shim的分层设计实现同一套 Zeppelin 支持多个 Spark 版本与多个 Scala 版本的并给出各模块的职责、关键配置项与底层调用链。一、为什么需要模块化Spark 与 Scala 的版本兼容性挑战Zeppelin 的 Spark 支持面临两个天然的兼容性难题这也是spark目录采用多模块拆分的第一原因Spark 与 Scala 强绑定Spark 各版本编译于特定 Scala 二进制版本2.12 / 2.13Spark 运行时依赖的 Scala 版本必须与 Zeppelin 中加载的 Scala 编译器一致Spark API 跨版本演进从 Spark 2.x 到 3.x、再到 4.xSparkSession、SQL 语法与 REPLSparkILoop接口均发生变化。因此 spark/pom.xml 将整个 Spark 支持拆分为五个 Maven 模块并在spark.version默认 3.5.8、spark.scala.version默认 2.12.20等属性上通过 profile 灵活切换。下面按 README 中的模块清单逐一展开。二、模块结构总览五大模块各司其职spark/README.md给出了如下模块划分对应仓库中的实际目录模块仓库路径职责interpreterspark/interpreter入口模块定义全部解释器类创建 SparkContext/SparkSession是 PySpark、IPySpark 等解释器的依赖基础spark-scala-parentspark/spark-scala-parent各 Scala 模块的公共父 POM统一声明 scala-compiler、scala-library、scala-reflect 等依赖scala-2.12spark/scala-2.12Scala 2.12 专用实现兼容 Spark 2.x 与 3.xscala-2.13spark/scala-2.13Scala 2.13 专用实现仅 Spark 3.x 及以上支持 2.13spark-commonspark/spark-common与 Scala 版本无关的公共工具SparkUtils从源码结构看interpreter模块中的 Java 类SparkInterpreter、PySparkInterpreter等不直接依赖具体 Scala 版本而是通过桥接类AbstractSparkScalaInterpreter与各 Scala 版本模块通信从而实现了Java 主逻辑一次编写、Scala 适配层按版本编译的目标。三、interpreter 入口模块SparkInterpreter 的动态加载机制SparkInterpreterspark/interpreter/src/main/java/org/apache/zeppelin/spark/SparkInterpreter.java是 Spark 支持的总入口它本身是 Java 实现关键职责如下。3.1 构造阶段读取 Zeppelin 专属属性构造方法会读取两个关键属性并设置scala.color系统属性// 是否在 REPL 输出中启用 Scala 语法着色默认 true if (Boolean.parseBoolean(properties.getProperty(zeppelin.spark.scala.color, true))) { System.setProperty(scala.color, true); } // 是否启用 Spark 版本支持性检查默认 true this.enableSupportedVersionCheck java.lang.Boolean.parseBoolean( properties.getProperty(zeppelin.spark.enableSupportedVersionCheck, true));同时维护一个 Scala 版本到实现类的映射表2.12 → SparkScala212Interpreter、2.13 → SparkScala213Interpreter。3.2 open()SparkConf 构建与旧属性兼容转换open()方法执行初始化值得注意的细节包括属性透传interpreter 配置中的所有非空属性都会写入SparkConf旧属性兼容zeppelin.spark.useHiveContext会转换为spark.useHiveContext传入 Sparkzeppelin.spark.concurrentSQLtrue会设置spark.scheduler.pool的调度模式为FAIRspark/interpreter/src/main/java/org/apache/zeppelin/spark/SparkStringConstants.java 中定义SCHEDULER_MODE_PROP_NAMEmaster 兜底逻辑若未配置spark.master依次回退到master属性、环境变量MASTER最后使用本地模式默认值DEFAULT_MASTER_VALUE即local[*]。3.3 按运行时 Scala 版本动态加载适配类这是整个多版本支持的核心机制。loadSparkScalaInterpreter(SparkConf)通过extractScalaVersion()决定运行时 Scala 版本优先读取zeppelin.spark.scala.version配置否则使用scala.util.Properties.versionString()检测当前类路径上的 Scala 版本。if (conf.contains(zeppelin.spark.scala.version)) { scalaVersionString conf.get(zeppelin.spark.scala.version); } else { scalaVersionString scala.util.Properties.versionString(); } // 2.12 / 2.13 之外的版本直接抛出 Unsupported scala version随后使用双检锁double-checked locking确保内部解释器类只被加载一次并优先从ZEPPELIN_HOME/interpreter/spark/scala-版本目录通过URLClassLoader加载对应实现类在 yarn-cluster 模式下ZEPPELIN_HOME不可用则直接使用当前 ClassLoader。这也是scala-2.12、scala-2.13两个模块产物必须被放置到该目录的原因见 spark/spark-scala-parent/pom.xml 中maven-jar-plugin将输出目录设置为../../interpreter/spark/scala-${spark.scala.binary.version}。3.4 版本检查、会话管理与委托open()完成后SparkInterpreter会通过SparkVersion.fromVersionString(sc.version())获取运行时 Spark 版本若zeppelin.spark.enableSupportedVersionChecktrue且版本不受支持则拒绝启动。解释执行、取消、补全、进度获取等操作全部委托给内部 Scala 解释器internalInterpret执行前设置sc.setJobGroup()任务分组与spark.scheduler.pool本地属性实现段落级任务隔离与调度池切换cancel通过cancelJobGroup取消对应段落的任务getProgress委托内部实现基于 SparkStatusTracker统计已完成 task 比例。四、Scala 版本适配模块scala-2.12 与 scala-2.13两个 Scala 模块都继承 Java 桥接类AbstractSparkScalaInterpreterspark/interpreter/src/main/java/org/apache/zeppelin/spark/AbstractSparkScalaInterpreter.java该桥接类负责两者共享的生命周期逻辑。4.1 共享的 SparkContext/SparkSession 创建流程AbstractSparkScalaInterpreter.createSparkContext()是整个 Spark 会话创建的核心其逻辑包括若配置了spark.sql.catalogImplementationhive或旧属性zeppelin.spark.useHiveContexttrue则尝试启用 Hive 支持——前提是类路径中同时存在hive-site.xml且能加载 Hive 类通过Class.forName(org.apache.spark.sql.hive.HiveSessionStateBuilder)等探测否则降级为普通会话将zeppelin.interpreter.localRepo目录下的 jar 通过sc.addFile()分发给 executor把sparkSparkSession、scSparkContext、sqlContextSQLContext、zZeppelinContext四个对象绑定到 REPL 命名空间并在初始化时静默执行一组预置 importimport org.apache.spark.SparkContext._ import spark.implicits._ import sqlContext.implicits._ import spark.sql import org.apache.spark.sql.functions._这解释了为什么用户直接在 Spark 段落里调用sql(...)、spark、sc、z即可工作通过sc.uiWebUrl()获取 Spark Web UI 地址并支持zeppelin.spark.uiWebUrl模板其中{{applicationId}}会被替换为实际 applicationId在 YARN 模式下若开启spark.webui.yarn.useProxytrue则通过YarnClient查询 ApplicationReport 获取代理 URL。4.2 REPL 输出与代码补全的处理细节两个实现类spark/scala-2.12/src/main/scala/org/apache/zeppelin/spark/SparkScala212Interpreter.scala 与 spark/scala-2.13/src/main/scala/org/apache/zeppelin/spark/SparkScala213Interpreter.scala在 REPL 输出处理上保持一致使用InterpreterOutputStream捕获 Scala REPL 的 stdout将其重定向到段落输出并可通过printREPLOutput本地属性关闭当代码返回IR.Error且错误信息包含value toDF is not a member of ...RDD时自动重试并在代码前补上import sqlContext.implicits._规避 Scala 编译器 SI-6649 问题当返回IR.Incomplete时自动在代码末尾追加print()以消除以注释结尾导致的解析不完整处理 REPL 类名冲突open()中设置scala.repl.name.line为基于hashCode()的唯一前缀负数转成0确保多个 Scala REPL 并发时生成的临时类不会互相覆盖spark/interpreter/src/main/java/org/apache/zeppelin/spark/AbstractSparkScalaInterpreter.java 中注释详述了 scoped 模式下的类名冲突问题。两者差异主要在 REPL 实现细节2.13 版本基于SparkILoopReplCompletion且代码补全时需要反射兼容 Scala 2.13.7 前后CompletionCandidate字段名变化namevsdefString见 ZEPPELIN-5946 修复注释。4.3 会话关闭的资源清理close()会执行YARN 模式下清理spark.yarn.stagingDir或 HDFS 家目录下的.sparkStaging/applicationId临时目录、停止 SparkContext 与 SparkSession、关闭 SparkILoop并清空sqlContext、z引用。五、spark-common跨版本的公共工具 SparkUtilsspark/spark-common/src/main/java/org/apache/zeppelin/spark/SparkUtils.java 是唯一与 Scala 版本无关的公共模块提供两类能力Spark Job 事件监听setupSparkListener()注册SparkListener在onJobStart时当spark.ui.enabledtrue且zeppelin.spark.ui.hiddenfalse构建 Spark Job 的 Web URL 并推送供前端展示段落关联的 Spark Job 链接DataFrame 展示showDataFrame()将Dataset转为 Zeppelin 的%table表格输出支持template本地属性渲染单行模板SingleRowInterpreterResult并严格截断到maxResult行多取一行用于判断是否超出上限。六、多版本支持的另一面SparkVersion 检测与 Maven Profiles6.1 运行时版本检测spark/interpreter/src/main/java/org/apache/zeppelin/spark/SparkVersion.java 提供版本解析与比较能力将sparkContext.version()返回的字符串解析为major.minor.patch编码成 5 位整数如 2.0.0 → 20000用于比较。当前仓库定义最低支持版本MIN_SUPPORTED_VERSION 3.3.0未来未支持阈值UNSUPPORTED_FUTURE_VERSION 4.1.0无法识别的版本字符串会被视为未来版本编码 99999同样触发不受支持判定。6.2 构建期版本切换spark/interpreter/pom.xml 提供了两类 profileSpark 版本 profilespark-3.53.5.8、spark-3.43.4.3默认激活、spark-3.33.3.4以及spark-4.04.0.0各自还同步调整protobuf.version、py4j.version、libthrift.version等配套依赖版本Scala 版本 profilespark-scala-2.132.13.16与spark-scala-2.122.12.20默认激活通过spark.scala.binary.version决定依赖坐标如spark-core_2.12/spark-core_2.13。6.3 PySpark 文件的构建期打包interpreter模块的构建还承担了一个特殊任务通过download-maven-plugin下载 Spark 源码包并在generate-resources阶段用maven-antrun-plugin抽取其中python/lib/py4j-版本-src.zip与python/pyspark目录输出到 spark/interpreter/pyspark 目录供 PySpark 解释器运行时使用测试阶段会通过PYTHONPATH环境变量引用这些文件。这意味着 PySpark 客户端文件随 Zeppelin 构建自动附带无需用户额外准备。七、基于 SparkInterpreter 的衍生解释器README 中提到其他解释器PySparkInterpreter、IPySparkInterpreter 等都依赖 SparkInterpreter这在源码中体现得十分清晰PySparkInterpreterspark/interpreter/src/main/java/org/apache/zeppelin/spark/PySparkInterpreter.java继承PythonInterpreter在open()中通过getInterpreterInTheSameSessionByClassName(SparkInterpreter.class)获取同一会话中的 SparkInterpreter并复用其 SparkContext 与 ZeppelinContext执行前会通过 py4j 桥调用sc.setJobGroup()与sc.setLocalProperty(spark.scheduler.pool, ...)保证 Python 段落的 Job 归属与调度池与 Scala 段落一致。Python 解释器可执行文件按spark.pyspark.driver.python spark.pyspark.python PYSPARK_DRIVER_PYTHON PYSPARK_PYTHON的优先级选择IPySparkInterpreterspark/interpreter/src/main/java/org/apache/zeppelin/spark/IPySparkInterpreter.java功能更完整但要求 Jupyter/IPython 前置其开关由zeppelin.pyspark.useIPython默认 true控制SparkSqlInterpreterspark/interpreter/src/main/java/org/apache/zeppelin/spark/SparkSqlInterpreter.java提供纯 SQL 段落使用SqlSplitter拆分多条 SQL通过zeppelin.spark.concurrentSQL配合zeppelin.spark.concurrentSQL.max默认 10切换并行/串行调度器段落级limit本地属性控制结果行数默认取zeppelin.spark.maxResult即 1000错误输出由zeppelin.spark.sql.stacktrace控制是否打印完整堆栈zeppelin.spark.sql.interpolation控制 SQL 中的字符串插值是否启用。八、测试与验证多版本支持的证据spark模块在 spark/interpreter/src/test/java/org/apache/zeppelin/spark 下提供了系统化测试覆盖各解释器核心行为SparkInterpreterTest验证解释器初始化、%spark段落执行与z上下文SparkSqlInterpreterTest验证 SQL 执行、limit 截断与结果输出PySparkInterpreterTest与PySparkInterpreterMatplotlibTest验证 PySpark 初始化、matplotlib 绘图集成后者默认被排除出常规单测需显式运行SparkVersionTest验证版本解析与新旧比较逻辑另有IPySparkInterpreterTest、SparkUtilsTest等。测试环境通过PYTHONPATH引用构建期打包的pyspark.zip与py4j-版本-src.zip并设置ZEPPELIN_HOME指向仓库根目录surefire 配置-Xmx3072m保证 Spark 测试的内存需求。九、实战小结如何选择与配置 Spark 解释器在实际部署中理解上述分层结构有助于快速定位问题与调整配置选择 Spark 版本在构建时使用-Pspark-3.5、-Pspark-3.4、-Pspark-3.3或-Pspark-4.0profile默认 3.4.3选择 Scala 版本-Pspark-scala-2.12默认或-Pspark-scala-2.13产物会自动输出到interpreter/spark/scala-版本目录关键配置spark.master缺省时回退MASTER环境变量最后为local[*]、zeppelin.spark.maxResult结果行数上限默认 1000、zeppelin.spark.useHiveContextHive 支持、zeppelin.spark.concurrentSQL并行 SQL、zeppelin.spark.enableSupportedVersionCheck版本白名单检查默认开启等版本兼容前提当前仓库支持 Spark 3.3 至 4.0 范围低于 3.3.0 或高于等于 4.1.0 会被判定为不支持可通过关闭zeppelin.spark.enableSupportedVersionCheck尝试运行但不受官方保障。综上Spark interpreter 通过Java 入口 Scala 版本适配层 公共工具的模块化设计将版本差异收敛到scala-2.12/scala-2.13两个薄适配层与 Maven profile 中既保证了同一份核心逻辑会话管理、Job 追踪、结果渲染的稳定又最大化了对 Spark/Scala 版本矩阵的覆盖能力——这正是 spark/README.md 所述模块结构背后的真实工程意图。赞分享数据分析数据可视化大数据后端前端任务调度【免费下载链接】zeppelinWeb-based notebook that enables>项目地址https://gitcode.com/gh_mirrors/zeppe/zeppelin点击查看免费下载相关推荐3步、10分钟做出一份可启动的黑苹果 OpenCore EFIOpCore Simplify 一键生成指南3步、10分钟做出一份可启动的黑苹果 OpenCore EFIOpCore Simplify 一键生成指南 OpCore Simplify 是一款免费开源的开发工具CLI突破版本壁垒ServerWrecker多版本JAR支持架构深度解析突破版本壁垒ServerWrecker多版本JAR支持架构深度解析 一、版本兼容困境Minecraft压力测试工具的技术痛点 你是否曾因Minecraft服游戏开发后端开发工具终极指南如何在Zeppelin中实现Spark多版本集成与Scala支持终极指南如何在Zeppelin中实现Spark多版本集成与Scala支持 Apache Zeppelin作为一款强大的Web笔记本工具为数据科学家和工程师提数据分析数据可视化大数据后端上一篇如何在生产环境中部署nfnet_l0.ra2_in1kDocker容器化与API服务搭建下一篇无需越狱用Cowabunga Lite打造专属iOS界面零基础也能轻松上手创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表