资讯动态

Apache Spark Connect 概览:解耦式客户端-服务端架构的完整实战指南

发布时间:2026/9/19 22:35:32 来源:尧图企业网站定制
Apache Spark Connect 概览解耦式客户端-服务端架构的完整实战指南【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark导读本文以 Apache Spark 官方文档《Spark Connect Overview》为骨架系统讲解 Spark 3.4 引入的 Spark Connect 解耦式客户端-服务端架构它如何用「未解析逻辑计划 Protocol Buffers gRPC Arrow」把轻量客户端与 Spark 引擎彻底分离以及如何在 PySpark 与 Scala 中快速搭建 Connect 服务器、配置客户端连接、在独立应用中使用 Spark Connect并处理多租户稳定性、认证、流式 API 支持与反向代理路由等落地问题。读完本文你将掌握从零启动 Spark Connect Server、用sc://连接串在交互式 Shell 与独立程序中开发以及让 Spark Connect 安全接入 Kubernetes Ingress 的完整技术方案。一、什么是 Spark ConnectSpark Connect 是 Apache Spark 3.4 中引入的客户端-服务端解耦架构它允许通过DataFrame API 和未解析逻辑计划unresolved logical plans作为协议实现到 Spark 集群的远程连接。客户端与服务端的分离使得 Spark 及其开放生态可以从任何地方被调用——它可以被嵌入现代数据应用、IDE、Notebook 以及各类编程语言之中。其核心思想是客户端库不再包含完整的 Spark 引擎而是一个可以随处嵌入的薄 APIthin API。Spark Connect API 构建在 Spark 的 DataFrame API 之上以未解析逻辑计划作为客户端与 Spark Driver 之间**与语言无关language-agnostic**的协议。官方快速入门入口参见 Quickstart: Spark Connect该页面位于仓库docs/api/python/文档体系内。Spark Connect API 架构图上图展示了 Spark Connect 的 API 架构左侧是多样化的客户端入口现代数据应用、IDE/Notebook、多语言 SDK右侧是服务端 Spark Driver 内部的多租户应用网关Multi-tenant Application Gateway、分析器Analyzer、优化器Optimizer、调度器Scheduler与分布式执行引擎Distributed Execution Engine。客户端只负责调用 API全部计算逻辑在服务端完成。二、Spark Connect 的工作原理Spark Connect 客户端库的设计目标是简化 Spark 应用程序开发。其工作流程如下计划构建Spark Connect 客户端把 DataFrame 操作翻译为未解析的逻辑查询计划unresolved logical query plan。协议编码这些计划使用Protocol Buffersprotobuf进行编码。传输编码后的计划通过gRPC 框架发送到服务器。服务端解析Spark Connect 端点在 Spark Server 上接收并将未解析的逻辑计划翻译为 Spark 的逻辑计划算子。这与解析一条 SQL 查询类似——属性attributes和关系relations被解析构建出初始解析计划initial parse plan。标准执行随后标准 Spark 执行流程启动从而保证 Spark Connect 能继承 Spark 全部优化与增强能力Analyzer → Optimizer → Scheduler → Distributed Execution Engine。结果回传执行结果通过 gRPC 以Apache Arrow 编码的行批次row batches流式返回给客户端。Spark Connect 通信流程图从仓库源码看服务端入口由sbin/start-connect-server.sh启动其核心类为org.apache.spark.sql.connect.service.SparkConnectServer参见 start-connect-server.sh而请求的实际处理逻辑位于 SparkConnectService.scala它实现了 protobuf 定义的SparkConnectServicegRPC 服务如ExecutePlan等 RPC 方法。2.1 Spark Connect 客户端应用与传统 Spark 应用的区别Spark Connect 的核心设计目标之一是实现客户端与服务端的完全分离与隔离。因此开发者需要注意以下几点变化客户端不与 Spark Driver 运行在同一进程中客户端无法直接访问 Driver JVM 来操纵执行环境。特别是在 PySpark 中客户端不再使用 Py4J因此无法访问持有 DataFrame、Column、SparkSession 等 JVM 实现私有字段例如df._jdf。协议基于逻辑计划、不支持全部执行 APISpark Connect 协议以 Spark 逻辑计划作为抽象用声明式方式描述要执行的运算因此它不支持 Spark 的全部执行 API最重要的是不支持 RDD。基于会话session-based的客户端客户端无法访问会操纵所有已连接客户端共享集群环境的属性。最重要的是客户端无法访问静态 Spark 配置或 SparkContext。三、Spark Connect 的运维收益新架构为多租户环境缓解了以下经典运维问题稳定性Stability内存使用过高的应用只会影响自身环境因为它们可以在各自独立的进程中运行。用户可以在客户端自行定义依赖无需担心与 Spark Driver 产生依赖冲突。可升级性UpgradabilitySpark Driver 可以独立于应用无缝升级例如受益于性能改进与安全补丁。只要服务端 RPC 定义设计为向后兼容应用就能保持前向兼容forward-compatible。可调试性与可观测性Debuggability and observabilitySpark Connect 支持在开发阶段直接从你喜欢的 IDE 中进行交互式调试同时应用可以使用其所在框架原生的指标metrics与日志库进行监控。四、如何启动 Spark Server 并启用 Spark ConnectSpark Connect 目前支持PySpark 与 Scala两类应用。下面演示如何运行一个带 Spark Connect 的 Spark Server并让客户端应用通过 Spark Connect 客户端库连接它。4.1 下载并启动带 Spark Connect 的 Spark Server首先从 Apache Spark 官方下载页面获取 Spark 发行包选择最新版本以及适合的包类型通常选择 Pre-built for Apache Hadoop 3.5 and later。下载后解压tar -xvf spark-3.5.x-bin-hadoop3.tgz打开终端进入解压后的spark目录运行start-connect-server.sh脚本启动带 Spark Connect 的 Spark Server./sbin/start-connect-server.sh说明请确保所用包版本与下载的 Spark 版本一致。仓库中 start-connect-server.sh 的脚本逻辑是通过spark-daemon.sh submit org.apache.spark.sql.connect.service.SparkConnectServer以后台守护进程方式提交 Connect 服务器且支持--wait参数以前台方式运行便于调试。除独立服务器外还可以使用./bin/spark-connect-shell启动一个交互式 Scala Shell其Connect 服务器内嵌在 Shell 进程内部适合本地快速体验。启动完成后Spark Server 即处于运行状态可以接受来自客户端应用的 Spark Connect 会话。4.2 用于交互式分析的 Spark Connect创建 Spark 会话时可以通过多种方式指定使用 Spark Connect。若未使用以下任一机制Spark 会话仍将和以前一样工作即不启用 Spark Connect。方式一设置SPARK_REMOTE环境变量在客户端机器上设置SPARK_REMOTE环境变量后创建新 Spark 会话该会话即为 Spark Connect 会话。这种方式无需改动任何代码。export SPARK_REMOTEsc://localhost ./bin/pysparkPySpark Shell 启动后欢迎信息会提示已通过 Spark Connect 连接Client connected to the Spark Connect server at localhost方式二创建 Spark 会话时显式指定也可以在启动 PySpark Shell 时通过remote参数指定服务器位置./bin/pyspark --remote sc://localhost同样可以看到欢迎信息提示已连接。还可以检查会话类型来确认——如果类型路径包含.connect.即表示在使用 Spark ConnectSparkSession available as spark. type(spark) class pyspark.sql.connect.session.SparkSession然后即可正常运行 PySpark 代码验证 Spark Connect 生效 columns [id, name] data [(1,Sarah), (2,Maria)] df spark.createDataFrame(data).toDF(*columns) df.show() -------- | id| name| -------- | 1|Sarah| | 2|Maria| --------Scala Shell 方式Scala Shell 基于 Ammonite REPL用法与 PySpark Shell 类似./bin/spark-shell --remote sc://localhostREPL 初始化成功后会显示欢迎横幅并默认尝试连接本机 Spark ServerWelcome to ____ __ / __/__ ___ _____/ /__ _\ \/ _ \/ _ / __/ _/ /___/ .__/\_,_/_/ /_/\_\ version 3.5.x /_/ Type in expressions to have them evaluated. Spark session available as spark.运行 Scala 代码验证 spark.range(10).count res0: Long 10L4.3 配置客户端-服务器连接默认情况下REPL 会尝试连接本机 15002 端口的 Spark Server该端口也是sc://连接串的默认端口。连接可以通过以下几种方式配置完整参数规范见仓库 client-connection-string.md。设置SPARK_REMOTE环境变量在 REPL 启动时自定义客户端-服务器连接export SPARK_REMOTEsc://myhost.com:443/;tokenABCDEFG ./bin/spark-shell或直接内联SPARK_REMOTEsc://myhost.com:443/;tokenABCDEFG spark-connect-repl通过连接字符串以编程方式创建连接使用SparkSession#builder import org.apache.spark.sql.SparkSession val spark SparkSession.builder.remote(sc://localhost:443/;tokenABCDEFG).getOrCreate()连接字符串参数详解sc://连接串遵循标准 URI 定义scheme 固定为sc://路径组件必须为空参数通过 HTTP URL Path Parameter 语法传递且所有参数区分大小写sc://host:port/;param1value;param2value参数类型说明示例hostStringSpark Connect 端点主机名必须为完整域名或 IP 地址gRPC 端点不支持路径myexample.com、127.0.0.1portNumericgRPC 端点端口默认 1500215002、443tokenString设置后启用基于标准 bearer token 的 gRPC 认证设置该值会同时启用 SSLtokenABCDEFGHuse_sslBoolean使用 TLS 连接默认false证书需在系统信任库中use_ssltrueuser_idString自动写入 Spark Connect UserContext 消息的用户 ID用于 Spark 会话管理可选user_idMartinuser_agentString代表用户执行请求的应用标识Python 客户端默认_SPARK_CONNECT_PYTHONuser_agentmy_data_query_appsession_idString服务端会话缓存的键可用于跨语言共享同一 Spark 会话须为合法 UUID默认随机生成session_id550e8400-e29b-41d4-a716-446655440000grpc_max_message_sizeNumericgRPC 消息最大字节数默认128 * 1024 * 1024grpc_max_message_size134217728grpc_keepalive_enabledBoolean是否发送 gRPC/HTTP2 keepalive PING 以探测静默死连接如 NAT 网关/负载均衡丢弃空闲连接默认truegrpc_keepalive_enabledfalsegrpc_keepalive_time_msNumeric空闲多少毫秒后发送 keepalive PING。服务端容忍客户端 PING 频率不低于每 10s 一次低于该下限会被以too_many_pings断开grpc_keepalive_time_ms30000grpc_keepalive_timeout_msNumeric等待 PING 确认的毫秒数超时判定连接死亡grpc_keepalive_timeout_ms10000grpc_keepalive_without_callsBoolean无在途 RPC 时是否仍持续发送 keepalive PING默认truegrpc_keepalive_without_callsfalse典型用法示例server_url sc://myhost.com/ # 默认端口 15002 server_url sc://myhost.com:443/;use_ssltrue # 换端口 TLS server_url sc://myhost.com:443/;use_ssltrue;tokenABCDEFG # TLS 令牌认证 server_url sc://myhost.com:443/;grpc_keepalive_time_ms30000;grpc_keepalive_timeout_ms10000 # 加速探测死连接需要特别注意的是连接串中不能携带路径前缀如sc://myhost.com:443/mypathprefix/;tokenAAAAAAA是非法用法因为 gRPC 的方法名本身就是 HTTP/2 的:path客户端无法指定独立路由路径这在下文「反向代理路由」一节还会详述。五、更快地本地迭代持久化 Connect 服务器在本地开发或测试时如果用from pyspark.sql import SparkSession spark SparkSession.builder.remote(local[*]).getOrCreate()PySpark 会启动一个仅在当前 Python 进程存活期间存在的全新进程内 Connect 服务器。每次执行python script.py或每个新的测试 worker 进程都要重新支付一次性启动成本——JVM 预热、SparkContext构建、Connect 服务器启动——这些往往需要数秒让「改代码 → 跑一下」的循环变得缓慢。5.1 手动启动持久化服务器为了摊销这一成本可以启动一个持久化的本地 Connect 服务器并让每次运行都连接它# 只启动一次跨多次运行保持在线。--master 可选默认 local[*] $SPARK_HOME/sbin/start-connect-server.sh --master local[*] # 每次运行都重连而不是重新启动一个新服务器 python -c from pyspark.sql import SparkSession; SparkSession.builder.remote(sc://localhost:15002).getOrCreate() # 用完停止 $SPARK_HOME/sbin/stop-connect-server.sh5.2 由 PySpark 托管的持久化服务器实验性在 POSIX 系统上PySpark 可以代为管理这个持久化服务器。设置SPARK_LOCAL_CONNECT_REUSE1或 builder 上的spark.local.connect.reusetrue后SparkSession.builder.remote(local[*]).getOrCreate()会在首次运行时通过sbin/start-connect-server.sh启动一个持久化服务器后续运行则直接重连脚本里可以一直使用普通的local[*]URLexport SPARK_LOCAL_CONNECT_REUSE1 # 第一次运行启动服务器后续运行重连 python -c from pyspark.sql import SparkSession; SparkSession.builder.remote(local[*]).getOrCreate() # 完成时停止托管服务器 python -m pyspark.sql.connect.local_server --stop从源码看该机制的实现在 local_server.py托管服务器是一个普通的spark-daemon.sh守护进程但它运行在按用户隔离的 pid 目录与标识串下因此不会与手工启动的服务器冲突——这也意味着普通的sbin/stop-connect-server.sh找不到它。--stop命令会向记录的服务器发送信号并清理发现文件直接 kill 服务器 pid 同样有效下一次运行发现服务器已死会自动启动新实例。连接细节host、port、认证 token、pid、Spark 版本记录在每个用户私有的发现文件中可用SPARK_LOCAL_CONNECT_DISCOVERY覆盖其位置。发现文件与日志存放于系统临时目录下按用户隔离的0700权限目录token 以0600权限存储且服务器总是绑定 IPv4 loopback覆盖任何配置的绑定地址因此机器上其他用户既读不到 token 也无法认证——同用户的不同进程按设计共享服务器。一次运行只重连Spark 版本匹配的服务器。升级 Spark 后旧版本服务器无法复用下一次运行会报错并提示停止旧服务器执行上面的--stop后重跑即可启动新服务器。会话隔离语义每次运行都建立独立的 Connect 会话因此会话级状态——临时视图、运行时 SQL 配置、会话产物——每次运行都是全新的绝不会在运行间泄漏而由共享SparkContext支撑的状态持久化 catalog/warehouse、全局临时视图、缓存的数据集会在运行间共享因此若要求运行间完全隔离请自行按运行命名空间划分数据库或清理这些状态。该托管工作流目前是实验性的--stop命令、发现文件位置与格式可能在未来版本变化例如本地服务器管理被并入统一的spark connectCLI。该机制依赖sbin/下的 POSIX 脚本Windows 不支持。六、在独立应用程序中使用 Spark Connect6.1 PythonPySpark独立应用首先安装客户端包。可以单独安装轻量客户端不含完整引擎pip install pyspark-client3.5.x如果是打包 PySpark 应用/库在setup.py中声明依赖install_requires[ pyspark-client3.5.x ]编写代码时创建 Spark 会话时通过remote函数引用你的 Spark Serverfrom pyspark.sql import SparkSession spark SparkSession.builder.remote(sc://localhost).getOrCreate()下面是一个完整的简单应用SimpleApp.py——统计一个文本文件中包含字母a与b的行数SimpleApp.py from pyspark.sql import SparkSession logFile YOUR_SPARK_HOME/README.md # Should be some file on your system spark SparkSession.builder.remote(sc://localhost).appName(SimpleApp).getOrCreate() logData spark.read.text(logFile).cache() numAs logData.filter(logData.value.contains(a)).count() numBs logData.filter(logData.value.contains(b)).count() print(Lines with a: %i, lines with b: %i % (numAs, numBs)) spark.stop()注意需要将YOUR_SPARK_HOME替换为 Spark 实际安装路径。用普通 Python 解释器运行$ python SimpleApp.py ... Lines with a: 72, lines with b: 396.2 Scala 独立应用在 Scala 应用/工程中使用 Spark Connect首先需要引入正确的依赖。以sbt为例在build.sbt中添加libraryDependencies org.apache.spark %% spark-connect-client-jvm % 3.5.x编写代码时同样通过remote函数创建会话import org.apache.spark.sql.SparkSession val spark SparkSession.builder().remote(sc://localhost).getOrCreate()关于用户自定义代码的重要说明涉及引用用户自定义代码的操作如 UDF、filter、map等需要注册一个 ClassFinder 来拾取并上传所需的 class 文件同时任何 JAR 依赖都必须通过SparkSession#addArtifact上传到服务器。示例import org.apache.spark.sql.connect.client.REPLClassDirMonitor // 注册一个 ClassFinder 来监视并上传构建输出目录中的 class 文件 val classFinder new REPLClassDirMonitor(ABSOLUTE_PATH_TO_BUILD_OUTPUT_DIR) spark.registerClassFinder(classFinder) // 上传 JAR 依赖 spark.addArtifact(ABSOLUTE_PATH_JAR_DEP)其中ABSOLUTE_PATH_TO_BUILD_OUTPUT_DIR是构建系统写出 class 文件的输出目录ABSOLUTE_PATH_JAR_DEP是本地文件系统上 JAR 的位置。REPLClassDirMonitor是ClassFinder的一个内置实现用于监视指定目录你也可以继承ClassFinder实现自定义的搜索与监视逻辑。关于 Spark Connect 应用开发与自定义扩展的更多内容参见仓库文档 Application Development with Spark Connect。七、客户端应用认证Spark Connect没有内置认证机制但它被设计为可以无缝对接你现有的认证基础设施。其 gRPC HTTP/2 接口支持使用认证代理authenticating proxies因此可以在不向 Spark 中直接实现认证逻辑的情况下保护 Spark Connect。八、支持范围What is supportedPySpark自 Spark 3.4 起Spark Connect 支持大部分 PySpark API包括 DataFrame、Functions 与 Column。但不支持SparkContext、RDD 等 API。可以在 PySpark 的 API reference 文档中逐项核对被标注为Supports Spark Connect的 API 即表示已支持迁移既有代码前可据此检查。Scala自 Spark 3.5 起Spark Connect 支持大部分 Scala API包括 Dataset、functions、Column、Catalog 与 KeyValueGroupedDataset。UDF用户自定义函数User-Defined Functions受支持——在 Shell 中默认支持在独立应用中则需要额外配置即上文提到的 ClassFinder 与 AddArtifact。流式 API大部分 Streaming API 受支持包括 DataStreamReader、DataStreamWriter、StreamingQuery 与 StreamingQueryListener。不支持SparkContext 与 RDD 在 Spark Connect 中均不受支持。更多 API 支持正在规划中预计在未来的 Spark 版本中推出。九、通过共享入口或反向代理路由Ingress/Reverse Proxy当多个服务在 Kubernetes Ingress或其他反向代理后面共享同一主机名时一个常见的想法是给每个服务一个 URL 路径前缀例如sc://host/sparkConnect并基于路径路由。这种做法对 gRPC 行不通gRPC 方法名本身就是 HTTP/2 的:path例如/spark.connect.SparkConnectService/ExecutePlan因此连接串无法携带独立的路由路径。给:path加前缀会产生未知方法服务器会返回UNIMPLEMENTED——除非代理被配置为在转发前把前缀剥掉而这是一个不属于 gRPC 设计范畴的双边契约。真正自由的路由维度是 HTTP/2 的:authority虚拟主机。做法是通过grpc.default_authoritychannel 选项把它设为一个路由标签routing tag然后在代理侧基于它路由例如 Ingress 的host:规则。客户端仍然拨号共享主机名default_authority只覆盖用于路由的:authority头不会改变实际连接的地址而 gRPC 方法的:path从未被改动因此无需任何路径重写。这正是 gRPC 维护者针对该场景推荐的做法。下面的示例使用 Python 客户端它可以直接暴露 gRPC channel 选项from pyspark.sql.connect.session import SparkSession from pyspark.sql.connect.client import DefaultChannelBuilder cb DefaultChannelBuilder(sc://myhost.com:443) cb.setChannelOption(grpc.default_authority, sparkconnect) # 路由标签 spark SparkSession.builder.channelBuilder(cb).getOrCreate()对应的 Kubernetes Ingress 基于该标签路由并让服务保持在/无子路径、无重写。注意路由标签必须是小写的 RFC 1123 名称因为 Kubernetes Ingress 的host:字段有此要求spec: rules: - host: sparkconnect # 匹配 grpc.default_authority http: paths: - path: / pathType: Prefix backend: service: name: spark-connect-server port: { number: 15002 }在 TLS 场景下gRPC 默认把:authority用作证书校验名称因此用路由标签覆盖它会导致校验失败标签不在服务器证书的 SAN 中。解决办法是让两个名称各司其职用grpc.ssl_target_name_override指定真实服务器主机名用于证书校验与 SNI用grpc.default_authority指定路由标签cb DefaultChannelBuilder(sc://myhost.com:443/;use_ssltrue) cb.setChannelOption(grpc.ssl_target_name_override, myhost.com) # 证书校验 / SNI cb.setChannelOption(grpc.default_authority, sparkconnect) # 路由标签 spark SparkSession.builder.channelBuilder(cb).getOrCreate()几点补充说明use_ssltrue会用系统可信 CA 库校验服务器证书。如果你的网关证书由客户端默认不信任的 CA 签发自签名或内部 CA应在客户端侧让该 CA 受信任例如通过GRPC_DEFAULT_SSL_ROOTS_FILE_PATH环境变量而不是通过连接串解决。代理自身的 TLS 配置对哪个主机名出示哪个证书属于 Ingress 配置范畴不在本文讨论范围。关于grpc.ssl_target_name_override的两个要点gRPC 将其文档化为面向测试的选项因为其典型误用是掩盖证书名称不匹配用一个服务器并未实际出示的名称去校验从而破坏主机名校验。而在上述场景中它被设置为真实、已校验的主机名证书仍被正确检查只是避免路由标签被用作校验名。如果你希望完全避开该选项可以给服务器证书的 SAN 中直接加上路由标签——此时仅设置grpc.default_authority即可无需任何 override。十、核心源码与文档索引如需进一步深入可在仓库中查阅以下关键路径服务端入口脚本sbin/start-connect-server.sh服务端 gRPC 实现SparkConnectService.scala连接串规范sql/connect/docs/client-connection-string.mdScala 客户端 ClassFinder 接口ClassFinder.scalaPython 本地持久化服务器实现python/pyspark/sql/connect/local_server.pySpark Connect 应用开发进阶指南docs/app-dev-spark-connect.md适用前提说明本文中的交互示例SPARK_REMOTE、--remote参数、pyspark/spark-shell与独立应用示例均要求先按「四、」一节启动对应版本的 Spark Connect 服务器且客户端与服务端使用相同的 Spark 版本pyspark-client与spark-connect-client-jvm的版本号请替换为与你的服务器一致的实际版本。持久化服务器复用SPARK_LOCAL_CONNECT_REUSE目前为实验特性仅支持 POSIX 系统。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价