ARTICLE DETAIL

资讯详情

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

Apache Spark 应用开发实战:Spark Connect 客户端-服务端解耦架构、API 模式与 Server Library 扩展

Apache Spark 应用开发实战:Spark Connect 客户端-服务端解耦架构、API 模式与 Server Library 扩展 Apache Spark 应用开发实战Spark Connect 客户端-服务端解耦架构、API 模式与 Server Library 扩展【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/sparkApache Spark 自 3.4 起引入 Spark Connect通过客户端-服务端解耦架构让基于 DataFrame API 与未解析逻辑计划unresolved logical plan的协议取代传统 Driver 内嵌客户端实现远程、多语言、可嵌入的 Spark 应用开发。本文以仓库 docs/app-dev-spark-connect.md 为骨架结合 sql/connect 模块源码与 docs/configuration.md 配置文档系统讲解 Spark Client Applications 与 Spark Server Libraries 两大开发范式、spark.api.mode与spark.remote两种接入方式以及 Relation / Expression / Command 三大协议扩展点的完整落地流程。读完本文你将能够用一行配置把经典应用切换到 Connect 模式、把客户端应用与 Spark 服务端彻底解耦并亲手实现一个自定义表达式插件端到端发布给 PySpark 客户端使用。Spark Connect为什么需要解耦的客户端-服务端架构在 Spark 3.4 之前Spark 应用程序通常运行在 Driver JVM 内部客户端代码与执行引擎共享同一进程用户只能通过 SQL、JDBC 等有限通道远程使用 Spark。Spark Connect 改变了这一点它引入客户端-服务端解耦架构允许通过DataFrame API和未解析逻辑计划作为协议从任意位置远程连接 Spark 集群。这种分离带来的直接效果是Spark 及其开放生态可以在任何地方被使用——可以嵌入现代数据应用、IDE、Notebook也可以嵌入任意编程语言Spark Connect 客户端目前官方支持 PySpark 与 Scala。协议本身是完全声明式的客户端只负责构建逻辑计划真正解析、优化与执行都在服务端完成。关于 Spark Connect 的整体概念、下载与启动方式、交互式分析与独立应用接入可继续阅读 docs/spark-connect-overview.md本文聚焦应用开发层面的两个角色划分与扩展机制。重新定义 Spark 应用Client Applications 与 Server LibrariesSpark Connect 将 Spark 应用开发者清晰地划分为两类角色二者围绕 Spark 的定位不同Spark 客户端应用Spark Client Applications使用 Spark 及其丰富生态进行分布式数据处理的常规应用例如 ETL 流水线、数据准备、模型训练与推理。它们通过 Spark Connect API本质就是 DataFrame API连接 Spark。Spark 服务器库Spark Server Libraries构建在 Spark 之上、扩展并补全 Spark 功能的库例如 MLlib使用 Spark 分布式算力的分布式机器学习库。Spark Connect 可以被扩展为 Server Libraries 暴露客户端侧的接口。图 1Spark Connect 架构示意——客户端与服务端通过完全声明式的 Spark Connect APIDataFrame API通信服务端以扩展点方式挂载自定义逻辑。如上图所示客户端应用通过Spark Connect API即 DataFrame API完全声明式连接 SparkSpark Server Libraries 则在服务端提供额外的服务端逻辑并通过Spark Connect 扩展点将其作为 Spark Connect API 的一部分暴露给客户端。例如图中的蓝色框Custom Library Plugin代表自定义服务端逻辑客户端通过 Spark Connect API 的蓝色部分调用它可以与 PySpark 或 Spark Scala 客户端并列使用让客户端应用轻松使用自定义库。Spark 3.4 及以后的演进目标是简化 Spark Client Applications 的开发同时为构建 Spark Server Libraries 提供清晰的扩展点与指南让两类应用都能与 Spark 一同平滑演进。Spark API 模式Spark Client 与 Spark Classic 的无缝切换Spark 提供了 API 模式配置项spark.api.mode让 Spark Classic 应用可以无缝切换到 Spark Connect。根据该配置的值应用可以运行在 Spark Classic 或 Spark Connect 两种模式之一。通过配置开启 Connect 模式Python 中构建会话时显式指定from pyspark.sql import SparkSession SparkSession.builder.config(spark.api.mode, connect).master(...).getOrCreate()也可以在提交 Scala 或 PySpark 应用时通过命令行配置spark-submit --master ... --conf spark.api.modeconnect根据 docs/configuration.md 中的官方配置说明spark.api.mode的默认值为classic可取值classic或connect自 4.0.0 版本引入配置项默认值含义引入版本spark.api.modeclassic对于 Spark Classic 应用指定是否通过启动本地 Spark Connect 服务器自动使用 Spark Connect取值为classic或connect4.0.0从源码看PySpark 在构建会话时解析该配置的逻辑位于 python/pyspark/sql/session.py优先读取opts中的spark.api.mode其次读取环境变量SPARK_API_MODE值不在(classic, connect)内时回退到 python/pyspark/util.py 中的default_api_mode()——它依据pyspark_connect模块是否可导入来决定默认模式。一旦判定为 connect 模式或检测到SPARK_REMOTE/spark.remote会话构建便会转入pyspark.sql.connect.session的RemoteSparkSession分支。本地测试的快捷方式spark.remoteSpark Connect 还提供了方便的本地测试选项。将spark.remote设置为local[...]或local-cluster[...]即可启动一个本地 Spark Connect 服务器并获得 Spark Connect 会话from pyspark.sql import SparkSession SparkSession.builder.remote(local[*]).getOrCreate()这与--conf spark.api.modeconnect搭配--master ...的效果类似但有两点关键区别spark.remote与--remote仅限于local*取值--conf spark.api.modeconnect搭配--master ...还支持Spark Classic 的更多集群 URL如spark://兼容性更广。spark.remote的解析逻辑可参见 sql/connect/common/src/main/scala/org/apache/spark/sql/connect/SparkSession.scala 的withLocalConnectServer它依次从 spark 配置、系统属性spark-submit注入、环境变量SPARK_REMOTE读取远程地址当远程地址以local开头或在 API 模式为 connect 时存在 master且发现连接启动脚本存在时就会拉起一个本地 Spark Connect 服务器进程并通过sc://localhost/;token...建立带认证令牌的连接进程退出时自动执行stop-connect-server清理。本地连接的三条常用路径详见 docs/spark-connect-overview.md设置环境变量SPARK_REMOTEsc://localhost后直接运行./bin/pyspark无需改动任何代码即进入 Connect 会话使用./sbin/start-connect-server.sh启动独立服务端默认绑定端口 15002见 Connect.scala 与配置文档客户端用SparkSession.builder.remote(sc://localhost:15002)连接在 Python 进程内使用SparkSession.builder.remote(local[*])PySpark 会在进程内启动一次性的 Connect 服务器如需跨进程复用可先启动持久化服务器$SPARK_HOME/sbin/start-connect-server.sh --master local[*]再让每次运行重连。Spark 客户端应用Spark Client ApplicationsSpark 客户端应用就是如今 Spark 用户开发的常规应用——ETL 流水线、数据准备、模型训练或推理通常基于声明式的 DataFrame / Dataset API 构建。使用 Spark Connect 后核心行为保持不变但有两点重要差异低层、非声明式的 APIRDD不能再被客户端应用直接使用。缺失的 RDD 功能以更高级别的 DataFrame API 替代方案提供。客户端应用不再直接访问 Spark Driver JVM与服务器完全分离。基于 Spark Connect 的客户端应用可以采用与以往任何作业相同的方式提交。与使用早期 Spark 版本3.4 及以下的经典应用相比它带来以下优势可升级性UpgradabilitySpark Connect API 抽象了服务端的变更与改进客户端与服务器 API 干净分离升级到新 Spark Server 版本无缝进行。简洁性Simplicity暴露给用户的 API 数量从 3 个减少到 2 个Spark Connect API 完全声明式对熟悉 SQL 的新用户非常容易学习。稳定性Stability客户端应用不再运行在 Spark Driver 上既不会引发服务端不稳定也不会受服务端不稳定影响。远程连通性Remote connectivity解耦架构让远程使用 Spark 不再局限于 SQL 与 JDBC任何应用都可以交互式地把 Spark 当作服务来用。向后兼容Backwards compatibility除 RDD 用法外Spark Connect API 与早期 Spark 版本代码兼容RDD 的替代 API 列表由 Spark Connect 提供。独立应用中如何连接补充在独立 Python 应用中安装pyspark-client后在创建 SparkSession 时通过remote指定服务器地址即可from pyspark.sql import SparkSession spark SparkSession.builder.remote(sc://localhost).appName(SimpleApp).getOrCreate() logData spark.read.text(YOUR_SPARK_HOME/README.md).cache() print(Lines with a: %i, lines with b: %i % ( logData.filter(logData.value.contains(a)).count(), logData.filter(logData.value.contains(b)).count())) spark.stop()Scala 独立应用则需在build.sbt中加入spark-connect-client-jvm依赖libraryDependencies org.apache.spark %% spark-connect-client-jvm % spark-versionimport org.apache.spark.sql.SparkSession val spark SparkSession.builder().remote(sc://localhost).getOrCreate()需要特别注意的是涉及用户自定义代码UDF、filter、map 等的操作在 Scala 客户端中要求注册ClassFinder以上传所需类文件JAR 依赖则须通过SparkSession#addArtifact上传到服务器——这是客户端与服务器解耦后代码从客户端到服务端的关键机制import org.apache.spark.sql.connect.client.REPLClassDirMonitor val classFinder new REPLClassDirMonitor(ABSOLUTE_PATH_TO_BUILD_OUTPUT_DIR) spark.registerClassFinder(classFinder) spark.addArtifact(ABSOLUTE_PATH_JAR_DEP)REPLClassDirMonitor是官方提供的ClassFinder实现用于监控构建输出目录并自动上传类文件你也可以实现自己的ClassFinder进行定制化搜索与监控。SPARK_REMOTE环境变量常量定义在 SparkConnectClient.scala其连接参数host、port、token、SSL、metadata 等封装在Configuration中见同文件Configuration样例类并支持可重连执行reattachable execute等高级行为。Spark 服务器库Spark Server Libraries在 Spark 3.4 之前对 Spark 的扩展例如 Spark ML、第三方 NLP 库都是像客户端应用一样构建与部署的。从 Spark 3.4 与 Spark Connect 开始Spark 提供了显式扩展点通过 Spark Server Libraries 扩展 Spark。这些扩展点把功能暴露给客户端与 Spark 中既有的扩展机制如SparkSessionExtensions、SparkPlugin有所不同——后两者属于服务端进程内的扩展而 Spark Connect 扩展点是跨客户端-服务端边界的。一个 Spark Server Library 由以下四部分构成Spark Connect 协议扩展下图蓝色框Proto API在协议层面新增自定义消息一个 Spark Connect 插件Plugin把自定义协议消息翻译为 Catalyst 逻辑计划/表达式扩展 Spark 的应用逻辑真正在服务端执行的业务逻辑客户端包把 Server Library 的应用逻辑暴露给 Spark 客户端应用与 PySpark 或 Scala Spark Client 并列使用。图 2Spark Server Library 的四个组成部分与标注步骤。(1) Spark Connect 协议扩展Relation、Expression 与 Command要扩展 Spark开发者可以扩展 Spark Connect 协议中的三类主要操作Relation关系/算子、Expression表达式与Command命令。这三类消息都在 oneof 联合类型中预留了google.protobuf.Any扩展字段message Relation { oneof rel_type { Read read 1; // ... google.protobuf.Any extension 998; } } message Expression { oneof expr_type { Literal literal 1; // ... google.protobuf.Any extension 999; } } message Command { oneof command_type { WriteCommand write_command 1; // ... google.protobuf.Any extension 999; } }这些扩展字段允许把任意 protobuf 消息作为 Spark Connect 协议的一部分进行序列化消息内容代表扩展实现的参数或状态。仓库中这三类协议的真实定义分别位于relations.protoRelation的oneof rel_type内预置Read/Project/Filter/Join/Sort/Limit/Aggregate/SQL等几十种内置关系扩展字段号为 998expressions.protoExpression的oneof expr_type扩展字段号为 999commands.protoCommand的oneof command_type扩展字段号为 999。此外 base.proto 中还存在多处repeated google.protobuf.Any extensions 999;的列表式扩展字段供更复杂的协议扩展场景使用。要构建一个自定义表达式类型开发者首先需要定义该表达式的自定义 protobuf 定义。例如定义一个带子表达式与自定义字段的表达式message ExamplePluginExpression { Expression child 1; string custom_field 2; }(2) Spark Connect 插件实现 (3) 自定义应用逻辑接下来开发者实现 Spark Connect 的ExpressionPlugin类基于 protobuf 消息的输入参数编写自定义应用逻辑class ExampleExpressionPlugin extends ExpressionPlugin { override def transform( relation: protobuf.Any, planner: SparkConnectPlanner): Option[Expression] { // Check if the serialized value of protobuf.Any matches the type // of our example expression. if (!relation.is(classOf[proto.ExamplePluginExpression])) { return None } val exp relation.unpack(classOf[proto.ExamplePluginExpression]) Some(Alias(planner.transformExpression( exp.getChild), exp.getCustomField)(explicitMetadata None)) } }插件接口在仓库中的真实定义是 ExpressionPlugin.javaOptionalExpression transform(byte[] relation, SparkConnectPlanner planner)。从接口注释可以看出两点实现约定插件类必须可无参构造、不应依赖内部状态每个已注册的扩展都会被传入Any实例若插件支持处理该类型则由其负责把对象构造成逻辑表达式并在必要时遍历其子节点。插件注册表 SparkConnectPluginRegistry.scala 维护了 relation / expression / command / getStatus 四类插件链既支持编译期通过relation[...]、expression[...]这类 builder 注入也支持运行时通过createConfiguredPlugins从 Spark 配置加载loadRelationPlugins/loadExpressionPlugins/loadCommandPlugins分别读取对应配置键。转换发生在 SparkConnectPlanner.scala 的transformExpressionPlugin它从注册表惰性遍历所有插件逐个调用transform取第一个返回非空结果的插件若没有任何插件认领该类型则抛出noHandlerFoundForExtension(typeUrl)异常。同理transformRelationPlugin第 265 行附近与handleCommandPlugin第 3309 行附近分别处理 Relation 与 Command 的扩展。打包与配置让 Spark 加载自定义逻辑应用逻辑开发完成后代码必须打包为JAR并配置 Spark 加载这些额外逻辑。相关的 Spark 配置项如下spark.jars指定包含自定义表达式应用逻辑的 JAR 文件位置spark.connect.extensions.expression.classes指定 Spark 加载的每个表达式扩展的完整类名。基于这些配置Spark 会在启动时加载对应值并使其可用于后续处理。官方配置文档docs/configuration.md中还列出了完整的扩展类配置家族它们都定义在 Connect.scala配置项默认值插件接口需实现的 trait引入版本spark.connect.extensions.relation.classes(none)org.apache.spark.sql.connect.plugin.RelationPlugin3.4.0spark.connect.extensions.expression.classes(none)org.apache.spark.sql.connect.plugin.ExpressionPlugin3.4.0spark.connect.extensions.command.classes(none)org.apache.spark.sql.connect.plugin.CommandPlugin3.4.0spark.connect.extensions.getStatus.classes(none)org.apache.spark.sql.connect.plugin.GetStatusPlugin4.1.0spark.connect.ml.backend.classes(none)org.apache.spark.sql.connect.plugin.MLBackendPlugin4.0.0这些配置均为逗号分隔的类名列表加载后由SparkConnectPluginRegistry.createConfiguredPlugins通过Utils.classForName反射实例化。仓库测试 SparkConnectPluginRegistrySuite.scala 提供了ExampleExpressionPlugin、ExampleRelationPlugin、ExampleCommandPlugin的完整参考实现与正反用例例如把配置类名指向不存在的类this.class.does.not.exist会正确抛出异常是理解插件写法的第一手资料。启动配置示例放在spark-defaults.conf或spark-submit --conf中spark.jars/path/to/example-library.jar spark.connect.extensions.expression.classescom.example.ExampleExpressionPlugin(4) Spark Server Library 客户端包服务端组件部署完成后任何客户端只要发送正确的 protobuf 消息即可使用它。以上述例子为例向 Spark Connect 端点发送如下消息负载即可触发扩展机制{ project: { input: { sql: { query: select * from samples.nyctaxi.trips } }, expressions: [ { extension: { typeUrl: type.googleapis.com/spark.connect.ExamplePluginExpression, value: \n\006\022\004\n\002id\022\006testval } } ] } }为了让该示例在 Python 中可用应用开发者需要提供一个把新表达式封装进 PySpark 的 Python 库。为任何表达式提供函数的最简单方式是接收一个 PySpark Column 实例作为参数返回一个应用了该表达式的新 Column 实例from pyspark.sql.connect.column import Expression import pyspark.sql.connect.proto as proto from myxample.proto import ExamplePluginExpression # Internal class that satisfies the interface by the Python client # of Spark Connect to generate the protobuf representation from # an instance of the expression. class ExampleExpression(Expression): def to_plan(self, session) - proto.Expression: fun proto.Expression() plugin ExamplePluginExpression() plugin.child.literal.long 10 plugin.custom_field example fun.extension.Pack(plugin) return fun # Defining the function to be used from the consumers. def example_expression(col: Column) - Column: return Column(ExampleExpression()) # Using the expression in the Spark Connect client code. df spark.read.table(samples.nyctaxi.trips) df.select(example_expression(df[fare_amount])).collect()这样客户端应用开发者在使用 PySpark 时就能像调用普通函数一样调用example_expression服务端则通过协议扩展、插件与应用逻辑完成真正的计算——客户端Python→ protobuf 协议 → 服务端插件 → Catalyst 表达式 → Spark 执行的完整链路就此打通。小结Spark Connect 把应用开发重构为两个清晰的角色客户端应用专注声明式 DataFrame API 与业务逻辑享受可升级性、简洁性、稳定性、远程连通与向后兼容五大收益服务器库通过 Relation / Expression / Command 三大协议扩展点、ExpressionPlugin插件接口与spark.connect.extensions.*配置把自定义逻辑发布给任意客户端。无论是用spark.api.modeconnect一键切换模式还是用spark.remotelocal[*]快速本地验证或是实现一个完整的 Server Library其核心机制都能在 sql/connect 模块的 proto 定义、插件注册表、planner 转换逻辑与配套测试中找到落点。建议按如下顺序继续深入先阅读 docs/spark-connect-overview.md 掌握服务器启动与交互式使用再对照 SparkConnectPluginRegistrySuite.scala 动手实现自己的第一个插件。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表