Spark Connect 开发者指南:连接字符串协议、Proto 消息演进与客户端代码生成
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/;param1value;param2value关键约束如下URI scheme 固定为sc://整体必须是一个合法 URI能被大多数系统正确解析例如主机名必须是合法主机名不能包含任意字符配置参数采用HTTP URL 路径参数Path Parameter语法传参与 JDBC 连接字符串风格类似路径组件path component必须为空所有参数均区分大小写case sensitive。1.3 连接参数完整参考表参数类型说明示例hostStringSpark Connect 端点的主机名。由于端点必须是完全 gRPC 兼容的端点不能指定特定路径主机名必须完全限定也可以是 IP 地址。myexample.com、127.0.0.1portNumeric连接 gRPC 端点时使用的端口默认值为15002可使用任何合法端口号。15002、443tokenString设置后启用标准 gRPC Bearer Token 认证默认不设置。设置该值会同时启用 SSL。tokenABCDEFGHuse_sslBoolean置为 true 时默认使用 TLS 连接端点前提是系统中有验证服务器证书所需的证书。默认值为false。use_ssltrue、use_sslfalseuser_idString自动写入 Spark ConnectUserContext消息中的用户 ID用于 Spark Session 的正确管理。可选参数在某些部署场景下可能通过其他方式自动注入。user_idMartinuser_agentString代表用户发起请求的客户端用户代理典型场景是使用 Spark Connect 实现功能、代表用户执行 Spark 请求的应用。Python 客户端默认值_SPARK_CONNECT_PYTHON。user_agentmy_data_query_appsession_idString除用户 ID 外Spark Connect 服务端的 Spark Session 缓存还以 session ID 作为缓存键。该参数允许显式提供 session ID例如实现同一用户跨语言共享 Spark Session。值必须是合法 UUID 字符串格式。默认值随机生成的 UUID。session_id550e8400-e29b-41d4-a716-446655440000grpc_max_message_sizeNumeric允许的 gRPC 消息最大字节数。默认值128 * 1024 * 1024即 134217728 字节。grpc_max_message_size134217728grpc_keepalive_enabledBoolean客户端是否发送 gRPC/HTTP2 keepalive PING 以探测静默死亡的连接例如 NAT 网关或负载均衡器丢弃空闲连接映射而未关闭 socket使阻塞调用报错而不是永远挂起。可作为逃生舱关闭例如在容易出现长时间停顿GC 暂停等的环境中避免误判断连。默认值true。grpc_keepalive_enabledfalsegrpc_keepalive_time_msNumeric客户端发送 keepalive PING 前的空闲时间毫秒。Spark Connect 服务端容忍客户端 PING 的频率不低于每 10 秒一次若设置低于该下限连接将因too_many_pings被断开。默认值60000。grpc_keepalive_time_ms30000grpc_keepalive_timeout_msNumeric客户端等待 keepalive PING 确认ack后判定连接死亡的时间毫秒。默认值20000。grpc_keepalive_timeout_ms10000grpc_keepalive_without_callsBoolean当连接上没有进行中的 RPC 时是否继续发送 keepalive PING。默认值true。grpc_keepalive_without_callsfalse1.4 有效与无效配置示例有效示例连接myhost.com的15002端口server_url sc://myhost.com/使用不同端口并启用 SSLserver_url sc://myhost.com:443/;use_ssltrue启用 SSL 并携带 Tokenserver_url sc://myhost.com:443/;use_ssltrue;tokenABCDEFG调优 gRPC keepalive例如比 60s/20s 默认值更快地探测死连接或完全关闭server_url sc://myhost.com:443/;grpc_keepalive_time_ms30000;grpc_keepalive_timeout_ms10000server_url sc://myhost.com:443/;grpc_keepalive_enabledfalse无效示例由于 Spark Connect 使用标准 gRPC 客户端为保持与 gRPC 标准及 HTTP 兼容服务端路径不可配置。以下写法无效server_url sc://myhost.com:443/mypathprefix/;tokenAAAAAAA1.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第二步安装 bufproto 代码生成工具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覆盖默认的官方二进制下载行为。默认构建则使用protocArtifactcom.google.protobuf:protoc:${protobuf.version}:exe:${os.detected.classifier}与pluginArtifactio.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),仅供参考

相关新闻

F16非线性六自由度飞机模型Simulink搭建与飞控验证实践

F16非线性六自由度飞机模型Simulink搭建与飞控验证实践

简介:这是一份面向航空工程、飞行控制与仿真技术学习者的F16战斗机非线性飞行动力学SIMULINK仿真资源,核心包含六自由度非线性模型、高/低保真气动系数数据、发动机模型和标准大气模型,适合用于飞行控制策略设计、飞行性能评估及故障诊断研究…

2026/9/21 15:19:00 阅读更多 →
Java解析698报文数据项:从字节流到TLV的工程实践

Java解析698报文数据项:从字节流到TLV的工程实践

简介:一份基于Java的698报文数据项解析示例代码,面向协议开发、报文调试以及有一定Java基础的工程师,帮助理解698报文中数据项的组织方式与逐项解析实现思路,适合作为入门参考或二次开发底稿。压缩包共94个文件,以81个…

2026/9/20 14:50:16 阅读更多 →
基于Ontology的自愈式企业操作系统:构建企业级AI的自主循环

基于Ontology的自愈式企业操作系统:构建企业级AI的自主循环

简介:聚焦Palantir Paragon 2025大会的深度解读PDF,面向企业数字化转型、AI落地与数据治理从业者,系统梳理AIP与Ontology技术如何驱动自愈式企业操作系统。内容覆盖CEO Alex Karp的价值创造哲学、FDE前线部署模式,以及建筑、医疗、…

2026/9/21 18:01:00 阅读更多 →

最新新闻

MATLAB晶粒生长模拟:蒙特卡洛Potts模型实现与应用

MATLAB晶粒生长模拟:蒙特卡洛Potts模型实现与应用

1. 项目背景与核心价值在材料科学研究领域,晶粒组织的演化过程直接影响着金属、陶瓷等材料的力学性能和物理特性。传统实验方法需要耗费大量时间和资源进行金相制备、热处理和显微观察,而计算机模拟技术为研究者提供了一种高效、低成本的替代方案。这个M…

2026/9/22 0:07:44 阅读更多 →
5年大厂面试官揭秘:奇拿面试题新手避坑指南

5年大厂面试官揭秘:奇拿面试题新手避坑指南

5年大厂面试官揭秘:奇拿面试题新手避坑指南 官方文档翻了三遍还是像看天书?别慌,这就是典型的【奇拿】场景。很多【新手避坑】指南只讲理论,却忽略了大厂面试官真正想听的那句人话。今天我就把底裤都扒了,带你用最短时间抓住【奇拿】考点的核心,让你下…

2026/9/22 0:07:44 阅读更多 →
qq恢复网站入门到精通:3步避坑,选型不踩雷

qq恢复网站入门到精通:3步避坑,选型不踩雷

qq恢复网站入门到精通:3步避坑,选型不踩雷 官方文档翻了三遍还是晕?别急,谁第一次看QQ找回账号的后台逻辑不是这样。官方流程太冗长,关键节点藏得深,导致你卡在“验证方式”和“数据同步”上,根本抓不住重点。今天咱们不念经,直接拆解从0到1搭…

2026/9/22 0:07:44 阅读更多 →
Excel VBA中Range.Value数组特性解析与应用

Excel VBA中Range.Value数组特性解析与应用

1. 深入理解VBA中Range.Value返回的数组特性在Excel VBA开发中,Range对象的Value属性是最基础也是最常用的功能之一。但许多开发者(包括我在早期)都曾在这个看似简单的操作上栽过跟头。今天我们就来彻底剖析这个日常操作背后的机制。关键发现…

2026/9/22 0:07:44 阅读更多 →
3个坑解决版本升级API全变:手写实现如何打广告核心逻辑

3个坑解决版本升级API全变:手写实现如何打广告核心逻辑

3个坑解决版本升级API全变:手写实现如何打广告核心逻辑 版本升级后 API 全变了,你写的代码直接报 AttributeError ,是不是瞬间血压飙升?别慌,这种时候硬啃新文档不如 手写实现 底层逻辑来得快。…

2026/9/22 0:07:44 阅读更多 →
ISO9001体系高频面试题:3年实战避坑指南与代码级解析

ISO9001体系高频面试题:3年实战避坑指南与代码级解析

ISO9001体系高频面试题:3年实战避坑指南与代码级解析 昨天刚带一个新人做审计,他手里拿着从网上复制的《质量手册》草稿,问我在“4.1…

2026/9/22 0:06:44 阅读更多 →

日新闻

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天 配置环境就卡半天?别怪机器慢,多半是你没选对工具链。在Java、Go或Python的项目现场, 手写实现…

2026/9/22 0:00:41 阅读更多 →
剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑 面试被问原理答不上来,是不是常态?别慌。很多开发者对着 GitHub 开源仓库里的代码发呆,看似简单实则暗藏玄机。今天这份【剑帝加点】速查手册,直接带你拆解核心实现,把面试必考的原理讲透。…

2026/9/22 0:00:41 阅读更多 →
手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优 复制来的代码跑不通不知道怎么调?别慌,这种“复制粘贴地狱”在开发圈太常见了。尤其是做 图片压缩网站…

2026/9/22 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/21 3:13:20 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/21 2:19:36 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/21 4:51:05 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/21 15:36:51 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/21 15:36:51 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/19 23:35:34 阅读更多 →