Spark Connect 开发者指南:连接字符串协议、Proto 消息演进与客户端代码生成
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
导读
本文基于 Apache Spark 仓库中sql/connect模块的开发者文档,系统讲解 Spark Connect 作为逻辑计划门面(logical plan facade)的实现机制,重点覆盖三大开发者主题:跨语言客户端统一遵循的sc://连接字符串协议(含全部参数的默认值与语义)、基于 proto3 的 Spark Connect 协议消息扩展规范,以及 Python 客户端代码生成与自定义protoc构建的完整实操流程。读完本文,你将掌握如何为 Spark Connect 贡献新客户端、如何向协议中安全新增消息字段,以及如何在受限编译环境中完成connect模块的构建与测试。
Spark Connect 是 Apache Spark 中实现**逻辑计划门面(logical plan facade)**的模块:客户端只负责构建逻辑计划,真正的执行由服务端 Spark 完成。该模块直接集成在 Spark 的构建体系中,其目录结构如下:
sql/connect/ ├── client/ # 客户端实现相关代码 ├── common/ # 协议公共部分(proto 定义、公共实现) ├── docs/ # 开发者文档(连接字符串、proto 消息扩展规范) ├── server/ # 服务端实现 └── shims/ # 版本兼容 shim需要注意的是,本模块文档面向 Spark Connect 的开发者而非最终用户,因此以下内容围绕协议设计约定与开发流程展开。
一、统一连接面:sc://连接字符串协议
1.1 设计背景
与 JDBC 或其他数据库连接类似,Spark Connect 采用**连接字符串(connection string)**承载连接端点所需的相关参数。从客户端视角看,Spark Connect 本质上就是一个普通的 gRPC 客户端,可以被标准 gRPC 方式配置;但为了让不同编程语言的客户端拥有一致的连接体验,Spark 在 sql/connect/docs/client-connection-string.md 中规范了统一的用户侧连接方式。
1.2 连接字符串语法
连接字符串遵循标准 URI 定义,其通用格式为:
sc://host:port/;param1=value;param2=value关键约束如下:
- URI scheme 固定为
sc://; - 整体必须是一个合法 URI,能被大多数系统正确解析,例如主机名必须是合法主机名,不能包含任意字符;
- 配置参数采用HTTP URL 路径参数(Path Parameter)语法传参,与 JDBC 连接字符串风格类似;
- 路径组件(path component)必须为空;
- 所有参数均区分大小写(case sensitive)。
1.3 连接参数完整参考表
| 参数 | 类型 | 说明 | 示例 |
|---|---|---|---|
host | String | Spark Connect 端点的主机名。由于端点必须是完全 gRPC 兼容的端点,不能指定特定路径;主机名必须完全限定,也可以是 IP 地址。 | myexample.com、127.0.0.1 |
port | Numeric | 连接 gRPC 端点时使用的端口,默认值为15002,可使用任何合法端口号。 | 15002、443 |
token | String | 设置后启用标准 gRPC Bearer Token 认证;默认不设置。设置该值会同时启用 SSL。 | token=ABCDEFGH |
use_ssl | Boolean | 置为 true 时默认使用 TLS 连接端点,前提是系统中有验证服务器证书所需的证书。默认值为false。 | use_ssl=true、use_ssl=false |
user_id | String | 自动写入 Spark ConnectUserContext消息中的用户 ID,用于 Spark Session 的正确管理。可选参数,在某些部署场景下可能通过其他方式自动注入。 | user_id=Martin |
user_agent | String | 代表用户发起请求的客户端用户代理,典型场景是使用 Spark Connect 实现功能、代表用户执行 Spark 请求的应用。Python 客户端默认值:_SPARK_CONNECT_PYTHON。 | user_agent=my_data_query_app |
session_id | String | 除用户 ID 外,Spark Connect 服务端的 Spark Session 缓存还以 session ID 作为缓存键。该参数允许显式提供 session ID,例如实现同一用户跨语言共享 Spark Session。值必须是合法 UUID 字符串格式。默认值:随机生成的 UUID。 | session_id=550e8400-e29b-41d4-a716-446655440000 |
grpc_max_message_size | Numeric | 允许的 gRPC 消息最大字节数。默认值:128 * 1024 * 1024(即 134217728 字节)。 | grpc_max_message_size=134217728 |
grpc_keepalive_enabled | Boolean | 客户端是否发送 gRPC/HTTP2 keepalive PING 以探测"静默死亡"的连接(例如 NAT 网关或负载均衡器丢弃空闲连接映射而未关闭 socket),使阻塞调用报错而不是永远挂起。可作为逃生舱关闭,例如在容易出现长时间停顿(GC 暂停等)的环境中避免误判断连。默认值:true。 | grpc_keepalive_enabled=false |
grpc_keepalive_time_ms | Numeric | 客户端发送 keepalive PING 前的空闲时间(毫秒)。Spark Connect 服务端容忍客户端 PING 的频率不低于每 10 秒一次;若设置低于该下限,连接将因too_many_pings被断开。默认值:60000。 | grpc_keepalive_time_ms=30000 |
grpc_keepalive_timeout_ms | Numeric | 客户端等待 keepalive PING 确认(ack)后判定连接死亡的时间(毫秒)。默认值:20000。 | grpc_keepalive_timeout_ms=10000 |
grpc_keepalive_without_calls | Boolean | 当连接上没有进行中的 RPC 时,是否继续发送 keepalive PING。默认值:true。 | grpc_keepalive_without_calls=false |
1.4 有效与无效配置示例
有效示例:连接myhost.com的15002端口:
server_url = "sc://myhost.com/"使用不同端口并启用 SSL:
server_url = "sc://myhost.com:443/;use_ssl=true"启用 SSL 并携带 Token:
server_url = "sc://myhost.com:443/;use_ssl=true;token=ABCDEFG"调优 gRPC keepalive,例如比 60s/20s 默认值更快地探测死连接,或完全关闭:
server_url = "sc://myhost.com:443/;grpc_keepalive_time_ms=30000;grpc_keepalive_timeout_ms=10000"server_url = "sc://myhost.com:443/;grpc_keepalive_enabled=false"无效示例:由于 Spark Connect 使用标准 gRPC 客户端,为保持与 gRPC 标准及 HTTP 兼容,服务端路径不可配置。以下写法无效:
server_url = "sc://myhost.com:443/mypathprefix/;token=AAAAAAA"1.5 源码视角:连接字符串如何被解析
连接字符串的解析与通道构建在 Python 客户端中由ChannelBuilder及其标准实现DefaultChannelBuilder完成,位于 python/pyspark/sql/connect/client/core.py:
- 参数常量与默认值:
ChannelBuilder定义了use_ssl、token、user_id、user_agent、session_id、grpc_keepalive_enabled、grpc_keepalive_time_ms、grpc_keepalive_timeout_ms、grpc_keepalive_without_calls等全部参数键(见PARAM_*常量),并定义了GRPC_MAX_MESSAGE_LENGTH_DEFAULT = 128 * 1024 * 1024。keepalive 相关默认值(GRPC_DEFAULT_KEEPALIVE_ENABLED = True、GRPC_DEFAULT_KEEPALIVE_TIME_MS = 60 * 1000、GRPC_DEFAULT_KEEPALIVE_TIMEOUT_MS = 20 * 1000、GRPC_DEFAULT_KEEPALIVE_WITHOUT_CALLS = True)与 JVM 客户端(SparkConnectClient.scala)保持一致(见代码中 SPARK-58094 注释)。 - scheme 校验:
DefaultChannelBuilder.__init__显式校验 URL 必须以sc://开头,否则抛出INVALID_CONNECT_URL错误;随后将sc://重写为http://以复用 Python 内置的urllib.parse,并校验path 组件必须为空。 - 参数解析:
_extract_attributes将参数段按;拆分、按=切分为键值对(非法格式会报错),并对值做 URL 解码;随后提取 hostname 与端口,未显式指定端口时使用DefaultChannelBuilder.default_port()(即 15002)。 - 安全通道决策:
secure属性定义为use_ssl or token is not None——即设置 token 会自动启用安全连接;当未启用 SSL 且主机为localhost时使用 gRPC 本地通道凭证(grpc.local_channel_credentials()),否则使用 SSL 通道凭证,token 通过grpc.access_token_call_credentials以组合凭证(composite credentials)方式附加。这也从实现层面印证了文档中"设置 token 会启用 SSL"的说明。
服务端则通过UserContext消息接收user_id,其定义在 sql/connect/common/src/main/protobuf/spark/connect/base.proto,包含user_id、user_name字段,并利用google.protobuf.Any类型支持扩展注入(repeated google.protobuf.Any extensions = 999)。AnalyzePlanRequest等请求消息中的session_id字段注释明确说明其格式应为 UUID 字符串(如00112233-4455-6677-8899-aabbccddeeff),由客户端设置以在同一 session 内汇总不同查询的流式响应。
二、协议演进规范:如何新增 Proto 消息与字段
Spark Connect 协议基于proto3定义,所有.proto文件位于 sql/connect/common/src/main/protobuf/spark/connect/,包含base.proto、relations.proto、expressions.proto、commands.proto、catalog.proto、types.proto、ml.proto、ml_common.proto、common.proto、example_plugins.proto、pipelines.proto等。由于 proto3 不再支持required约束,且非 message 类型的字段没有has_field_name函数来判断字段是否被设置,新增字段时需要遵循以下约定(详见 sql/connect/docs/adding-proto-messages.md)。
2.1 必填字段(Required)
新增具有必填语义的字段时,开发者必须遵循既定流程:对于服务端正确处理入站消息所必需的语义字段,必须在注释中以(Required)标注。对于标量字段(scalar fields),服务端不做额外的输入校验;对于复合字段(compound fields),服务端会做最小化检查以避免空指针异常,但不会做语义校验。
message DataSource { // (Required) Supported formats include: parquet, orc, text, json, parquet, csv, avro. string format = 1; }在 base.proto 中可以找到实际应用实例:AnalyzePlanRequest.session_id与AnalyzePlanRequest.user_context均以(Required)标注,表明它们是服务端处理该请求的必要输入。
2.2 可选字段(Optional)
语义上可选的字段必须使用optional关键字标记,服务端据此根据字段存在与否分支到不同的行为。由于标量类型缺乏可配置的默认值,可选值的单纯存在并不定义其默认值——服务端实现会根据自身规则解释观测到的值。
message DataSource { // (Optional) If not set, Spark will infer the schema. optional string schema = 2; }同样在 base.proto 中可以看到实例:client_observed_server_side_session_id是optional字段,服务端可用其校验服务端 session 是否已变化。需要留意的是,proto3 中的optional与 proto2 的optional语义不同:它显式追踪字段是否被设置(presence),这正是服务端分支判断的基础。
三、Python 客户端开发与代码生成
3.1 从 Proto 文件生成 Python 客户端代码
修改 Spark Connect 协议后,需要重新生成 Python 客户端代码。完整流程如下:
第一步:准备 Python 环境并安装依赖。具体要求是安装ruff以及 "Spark Connect python proto generation plugin (optional)" 一节中列出的依赖:
pip install --group dev第二步:安装 buf(proto 代码生成工具):
brew install bufbuild/buf/buf第三步:运行生成脚本:
dev/connect-gen-protos.sh该脚本位于 dev/connect-gen-protos.sh,其实现是对通用 proto 生成脚本的封装(./dev/gen-protos.sh connect "$@"),支持可选传入输出路径参数(./dev/connect-gen-protos.sh [path])。
3.2 生成产物的位置
生成的 Python 代码落入python/pyspark/sql/connect/proto/目录,包括base_pb2.py、base_pb2.pyi、base_pb2_grpc.py等文件。这些文件在仓库中已经存在,是运行上述生成脚本的产物,可以直接对照检查协议变更是否已同步到 Python 客户端。
四、自定义protoc与protoc-gen-grpc-java构建
4.1 适用场景
当编译环境中无法使用官方发布的protoc与protoc-gen-grpc-java二进制文件时(例如在默认glibc版本低于 2.14 的 CentOS 6 或 CentOS 7 上编译connect模块),可以通过指定用户自定义的protoc与protoc-gen-grpc-java二进制来编译和测试。
4.2 通过 Maven 构建
export SPARK_PROTOC_EXEC_PATH=/path-to-protoc-exe export CONNECT_PLUGIN_EXEC_PATH=/path-to-protoc-gen-grpc-java-exe ./build/mvn -Phive -Puser-defined-protoc clean package4.3 通过 sbt 构建
export SPARK_PROTOC_EXEC_PATH=/path-to-protoc-exe export CONNECT_PLUGIN_EXEC_PATH=/path-to-protoc-gen-grpc-java-exe ./build/sbt -Puser-defined-protoc clean package用户自定义的protoc与protoc-gen-grpc-java二进制可以在用户编译环境中通过源码编译产出,编译步骤参考 protobuf 与 grpc-java 官方构建说明。
4.4 构建配置的源码映射
上述 profile 在 sql/connect/common/pom.xml 中有明确的实现映射:Maven profileuser-defined-protoc将环境变量SPARK_PROTOC_EXEC_PATH映射为spark.protoc.executable.path、将CONNECT_PLUGIN_EXEC_PATH映射为connect.plugin.executable.path,并在protobuf-maven-plugin(版本 0.6.1)的配置中通过protocExecutable与pluginExecutable覆盖默认的官方二进制下载行为。默认构建则使用protocArtifact(com.google.protobuf:protoc:${protobuf.version}:exe:${os.detected.classifier})与pluginArtifact(io.grpc:protoc-gen-grpc-java:${io.grpc.version}:exe:${os.detected.classifier})自动下载与平台匹配的官方二进制。同样的user-defined-protocprofile 也存在于 common/config/pom.xml、core/pom.xml、sql/core/pom.xml、connector/protobuf/pom.xml 中,说明该机制适用于所有涉及 proto 代码生成的模块。
五、新客户端贡献指南
当为 Spark Connect 贡献新语言客户端时,需要意识到 Spark 致力于在所有语言间提供一致的用户体验,因此必须遵循以下两条核心指南:
- 连接字符串配置:严格遵循 sql/connect/docs/client-connection-string.md 中定义的
sc://连接字符串规范,保证不同语言客户端连接方式完全一致; - 新增协议消息:向 Spark Connect 协议新增消息时,必须遵守 sql/connect/docs/adding-proto-messages.md 中关于 proto3
(Required)/optional字段的标注约定,保证服务端跨语言行为统一。
这两份文档是协议层面的"单一事实来源"(single source of truth),任何语言客户端的实现都应与之一致,例如上文中 Python 客户端DefaultChannelBuilder的解析逻辑就是对该规范的直接落地实现。
总结
Spark Connect 的开发者生态围绕三个关键契约展开:以sc://连接字符串为核心的统一连接协议(含 token/SSL、session 管理、gRPC keepalive 调优等可配置参数);以 proto3 字段规则为基础的协议扩展约定((Required)与optional的语义边界);以及从 proto 定义到各语言客户端代码的生成与构建流水线(dev/connect-gen-protos.sh与user-defined-protocprofile)。理解这三层,即可在保持跨语言一致体验的前提下,安全地为 Spark Connect 贡献新客户端与协议能力。
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考