Apache Spark 应用开发实战:Spark Connect 客户端-服务端解耦架构、API 模式与 Server Library 扩展
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
Apache 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 Libraries
Spark 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 暴露客户端侧的接口。
图 1:Spark Connect 架构示意——客户端与服务端通过完全声明式的 Spark Connect API(DataFrame API)通信,服务端以扩展点方式挂载自定义逻辑。
如上图所示,客户端应用通过Spark Connect API(即 DataFrame API,完全声明式)连接 Spark;Spark 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.mode=connect根据 docs/configuration.md 中的官方配置说明,spark.api.mode的默认值为classic,可取值classic或connect,自 4.0.0 版本引入:
| 配置项 | 默认值 | 含义 | 引入版本 |
|---|---|---|---|
spark.api.mode | classic | 对于 Spark Classic 应用,指定是否通过启动本地 Spark Connect 服务器自动使用 Spark Connect;取值为classic或connect | 4.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.remote
Spark Connect 还提供了方便的本地测试选项。将spark.remote设置为local[...]或local-cluster[...],即可启动一个本地 Spark Connect 服务器并获得 Spark Connect 会话:
from pyspark.sql import SparkSession SparkSession.builder.remote("local[*]").getOrCreate()这与--conf spark.api.mode=connect搭配--master ...的效果类似,但有两点关键区别:
spark.remote与--remote仅限于local*取值;--conf spark.api.mode=connect搭配--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_REMOTE="sc://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 Applications)
Spark 客户端应用就是如今 Spark 用户开发的常规应用——ETL 流水线、数据准备、模型训练或推理,通常基于声明式的 DataFrame / Dataset API 构建。使用 Spark Connect 后,核心行为保持不变,但有两点重要差异:
- 低层、非声明式的 API(RDD)不能再被客户端应用直接使用。缺失的 RDD 功能以更高级别的 DataFrame API 替代方案提供。
- 客户端应用不再直接访问 Spark Driver JVM,与服务器完全分离。
基于 Spark Connect 的客户端应用可以采用与以往任何作业相同的方式提交。与使用早期 Spark 版本(3.4 及以下)的经典应用相比,它带来以下优势:
- 可升级性(Upgradability):Spark 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-version>"import 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 并列使用。
图 2:Spark 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.proto:
Relation的oneof rel_type内预置Read/Project/Filter/Join/Sort/Limit/Aggregate/SQL等几十种内置关系,扩展字段号为 998; - expressions.proto:
Expression的oneof expr_type,扩展字段号为 999; - commands.proto:
Command的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.java:Optional<Expression> 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.RelationPlugin | 3.4.0 |
spark.connect.extensions.expression.classes | (none) | org.apache.spark.sql.connect.plugin.ExpressionPlugin | 3.4.0 |
spark.connect.extensions.command.classes | (none) | org.apache.spark.sql.connect.plugin.CommandPlugin | 3.4.0 |
spark.connect.extensions.getStatus.classes | (none) | org.apache.spark.sql.connect.plugin.GetStatusPlugin | 4.1.0 |
spark.connect.ml.backend.classes | (none) | org.apache.spark.sql.connect.plugin.MLBackendPlugin | 4.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.classes=com.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.mode=connect一键切换模式,还是用spark.remote=local[*]快速本地验证,或是实现一个完整的 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),仅供参考