ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Spark 的 Spark Connect 客户端连接字符串(sc://)完整指南

Apache Spark 的 Spark Connect 客户端连接字符串(sc://)完整指南 大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载导读Spark Connect 为 Spark 引入了去中心化的客户端-服务器架构客户端Python、Scala、Java、R 等通过 gRPC 与远端的 Spark Connect Server 通信而无需在本地启动完整的 Spark 运行时。为了让不同语言的客户端拥有一致的连接入口Spark 定义了一套形如sc://host:port/;paramvalue的**连接字符串Connection String**规范。本文将基于仓库中的官方规范文档与 Python/JVM 客户端源码系统讲解连接字符串的语法、全部可用参数、SSL 与 Token 认证、gRPC keepalive 调优、多语言客户端中的实际用法以及从字符串到 gRPC Channel 的底层解析过程帮助你准确配置和排查 Spark Connect 客户端连接问题。背景为什么需要连接字符串从客户端角度看Spark Connect 本质上就是一个标准的 gRPC 客户端理论上可以直接按 gRPC 的方式配置端点。但为了让 Python、Scala、Java、R 等多种编程语言的使用体验保持一致并形成开箱即用的统一连接表面Spark 借鉴了 JDBC 等数据库连接的方式通过一个连接字符串来承载连接所需的全部参数。正如规范文档 sql/connect/docs/client-connection-string.md 开篇所述Spark Connect 使用一个包含相关参数的连接字符串这些参数被解析后用于连接 Spark Connect 端点。连接字符串通用语法连接字符串遵循标准的 URIRFC 2396定义sc://host:port/;param1value;param2value语法要点如下固定 schemeURI 的协议部分固定为sc://这是 Spark Connect 的专用标识。必须是合法 URI整个 URI 必须能被大多数系统正确解析。例如主机名必须合法不能包含任意字符。路径必须为空gRPC 端点不支持在 URL 中指定路径前缀因此sc://host:port后面的路径组件必须为空仅允许结尾的/。带路径前缀的写法是非法的见下文无效示例。参数采用 HTTP URL Path Parameter 风格所有配置参数以;分隔、以keyvalue形式放在路径之后这与 JDBC 连接字符串的写法风格相似。大小写敏感所有参数名和值都区分大小写。完整的 URI 形态为scheme://host:port/;key1value1;key2value2。注意参数区位于/之后、以;开头例如sc://myhost.com:443/;use_ssltrue。参数参考表下表完整列出 Spark Connect 连接字符串支持的全部参数来自 client-connection-string.md参数类型说明示例hostStringSpark Connect 端点的主机名。由于端点必须是完全 gRPC 兼容的端点因此无法指定特定路径。主机名必须是完全限定域名也可以是 IP 地址。myexample.com、127.0.0.1portNumeric连接 gRPC 端点时使用的端口。默认值15002。可以使用任何有效端口号。15002、443tokenString在 URL 中设置此参数时将启用使用 gRPC 的标准 Bearer Token 认证。默认不设置。设置该值会自动启用 SSL。tokenABCDEFGHuse_sslBoolean设置该标志后默认使用 TLS 连接端点。前提是系统中有用于验证服务器证书的必要证书。默认值false。use_ssltrue、use_sslfalseuser_idString自动设置到 Spark Connect UserContext 消息中的用户 ID用于 Spark Session 的合理管理。这是一个可选参数某些部署场景下可能通过其他方式自动注入。user_idMartinuser_agentString代表用户行事的客户端通常是基于 Spark Connect 实现功能、代表用户执行 Spark 请求的应用的用户代理标识。Python 客户端默认值_SPARK_CONNECT_PYTHON。user_agentmy_data_query_appsession_idStringSpark Connect Server 的 Session 缓存以 Session ID 作为缓存键。此参数允许显式提供 Session ID例如让同一用户跨多种语言共享 Spark Session。必须使用合法的 UUID v4 字符串格式。默认值随机生成的 UUID。session_id550e8400-e29b-41d4-a716-446655440000grpc_max_message_sizeNumericgRPC 消息允许的最大大小字节。默认值128 * 1024 * 1024128 MiB。grpc_max_message_size134217728grpc_keepalive_enabledBoolean客户端是否发送 gRPC/HTTP2 keepalive PING 来探测静默断开的连接例如 NAT 网关或负载均衡器丢弃空闲连接映射但未关闭 socket 的情况使阻塞调用以错误失败而不是永远挂起。在容易出现长停顿如 GC 暂停的环境中可关闭它以避免误判断连。默认值true。grpc_keepalive_enabledfalsegrpc_keepalive_time_msNumeric客户端发送 keepalive PING 前的空闲时间毫秒。Spark Connect Server 容忍客户端 PING 的频率不高于每 10 秒一次将此值设置低于该下限会被服务器以too_many_pings关闭连接。默认值60000。grpc_keepalive_time_ms30000grpc_keepalive_timeout_msNumeric客户端等待 keepalive PING 确认的时间毫秒超时即认为连接已死。默认值20000。grpc_keepalive_timeout_ms10000grpc_keepalive_without_callsBoolean当连接上没有进行中的 RPC 时是否继续发送 keepalive PING。默认值true。grpc_keepalive_without_callsfalse核心参数深入解析host 与 port端点的最小组成连接字符串中只有 host 是必填项端口不写时使用默认值15002。在 Python 客户端的 DefaultChannelBuilder 实现中默认端口常量同样定义为15002其default_port()方法在测试/开发模式下会通过 Py4J 从本地 Spark Session 获取实际启动端口生产环境一律回退到 15002。host 支持完全限定域名和 IPv4 地址若 host 是 IPv6 地址解析时会自动加上[]包裹见 core.py 中f[{self.url.hostname}]的处理。token 与 use_ssl认证与传输安全use_ssltrue表示使用 TLS 建立安全通道证书验证依赖系统信任库。token一旦设置隐式启用 SSLsecure属性为use_ssl or token is not None见 core.py即不能以明文方式携带 Token 传输。从源码实现看DefaultChannelBuilder.toChannel通道构建有三种分支非安全模式grpc.insecure_channel(endpoint)仅 Token 且 host 为localhost使用grpc.local_channel_credentials()本地凭据并复合access_token_call_credentials其他安全场景grpc.ssl_channel_credentials()复合 Bearer Token 凭据。Python 客户端还支持从环境变量SPARK_CONNECT_AUTHENTICATE_TOKEN读取 Token见 core.py便于在部署环境中避免把密钥写进 URL。user_id 与 session_idSession 隔离与共享Spark Connect Server 以(user_id, session_id)作为 Session 缓存键。user_id会写入 gRPC 请求的user_context字段见 core.py是 Session 归属与隔离的依据未显式指定时Python 客户端会回退到环境变量SPARK_USER再回退到USER见 core.py。session_id则用于跨语言/跨进程共享同一 Spark Session——只要user_id相同、session_id相同不同语言的客户端就能复用同一个服务端 Session。注意 core.py 会对session_id做 UUID v4 格式校验非法格式会抛出INVALID_SESSION_UUID_ID错误。user_agent客户端身份标识user_agent用于服务端日志与诊断标识代表用户执行请求的应用程序。Python 客户端默认值_SPARK_CONNECT_PYTHON且会在其后拼接spark/{版本} os/{系统} python/{版本}信息并对整个字符串做 URL 百分号编码若编码后超过 2048 字符会报错见 core.py。测试 test_client.py 验证了自定义user_agentbar会以^bar spark/... os/... python/...$的形式透传到服务端而默认值则为^_SPARK_CONNECT_PYTHON ...$。grpc_keepalive_*连接保活与死连接探测这组参数对应 SPARK-58094用于解决连接静默死亡但客户端永远挂起的问题NAT 网关或负载均衡器丢弃空闲连接映射但不发送 TCP RST/FIN 时阻塞的 RPC如流式查询的awaitTermination()会永远挂起开启 keepalive 后客户端周期性发送 HTTP/2 PING若收不到 PONG 确认则让阻塞调用以UNAVAILABLE错误返回。Python 客户端的默认值与 JVM 客户端保持一致grpc_keepalive_enabledtrue、grpc_keepalive_time_ms60000、grpc_keepalive_timeout_ms20000、grpc_keepalive_without_callstrue见 core.py。这些参数最终映射为 gRPC 的通道选项grpc.keepalive_time_ms、grpc.keepalive_timeout_ms、grpc.keepalive_permit_without_calls见 core.py。重要的服务器端约束Spark Connect Server 使用固定 10 秒的 PING 许可下限GRPC_KEEPALIVE_PERMIT_TIME_SECONDS因此grpc_keepalive_time_ms不要低于 10 秒否则会被服务器以too_many_pings强制断开。服务端集成测试 SparkConnectServiceKeepAliveSuite.scala 验证了 keepalive 端到端生效、健康长连接不受干扰以及关闭 keepalive 后调用会保持挂起等行为。有效连接示例以下示例来自规范文档的 Valid Examples 部分展示如何配置连接字符串。基础示例默认端口连接myhost.com的15002端口默认端口无需显式指定server_url sc://myhost.com/自定义端口 SSL连接myhost.com:443并使用 TLSserver_url sc://myhost.com:443/;use_ssltrueSSL Bearer Token 认证设置 Token 时自动启用 SSLserver_url sc://myhost.com:443/;use_ssltrue;tokenABCDEFGgRPC keepalive 调优将 keepalive 空闲时间从默认 60 秒缩短到 30 秒、确认超时从默认 20 秒缩短到 10 秒以便更快地探测到死连接注意 keepalive 间隔不应低于服务器 10 秒的许可下限server_url sc://myhost.com:443/;grpc_keepalive_time_ms30000;grpc_keepalive_timeout_ms10000或完全关闭 keepalive例如在容易出现长 GC 停顿、可能触发误判断连的环境server_url sc://myhost.com:443/;grpc_keepalive_enabledfalse多参数组合示例完整的生产环境连接串同时配置 SSL、Token、用户与 Session 共享、更大的消息上限server_url ( sc://myhost.com:443/;use_ssltrue;tokenABCDEFG; user_idMartin;session_id550e8400-e29b-41d4-a716-446655440000; grpc_max_message_size268435456 )无效示例由于 Spark Connect 使用标准的 gRPC 客户端为了与 gRPC 标准及 HTTP 保持兼容服务器路径不可配置。以下写法是非法的# 非法包含路径前缀 mypathprefix server_url sc://myhost.com:443/mypathprefix/;tokenAAAAAAA除路径必须为空外还有以下常见非法场景均会触发INVALID_CONNECT_URL错误参见 DefaultChannelBuilder.initscheme 不是sc://例如http://myhost.com:15002/缺少主机名例如sc:///;use_ssltrue参数不是keyvalue形式例如sc://host/;usessl缺少分隔符参数名或值大小写错误例如把use_ssl写成USE_SSL所有参数大小写敏感。多语言客户端中的使用方式PythonSparkSession.builder.remotePython 客户端PySpark通过SparkSession.builder.remote(url)传入连接字符串from pyspark.sql import SparkSession spark ( SparkSession.builder.remote(sc://myhost.com:443/;use_ssltrue;tokenABCDEFG) .build() )连接字符串由ChannelBuilder/DefaultChannelBuilder位于 python/pyspark/sql/connect/client/core.py解析先校验 scheme 必须为sc://再将sc://改写为http://以复用 Python 标准库的urllib.parse.urlparse进行 URI 解析随后从 URL 的 params 段按;拆分出keyvalue参数对并对值做 URL 解码见 core.py。Scala / JavaSparkConnectClient.BuilderJVM 客户端通过SparkConnectClient.builder()配置连接既可直接使用connectionString(url)也可使用与连接字符串参数一一对应的 Builder 方法host、port、token、useSsl、userId、userName、userAgent、sessionId、grpcMaxMessageSize、grpcKeepAlive* 等。命令行工具如 Connect REPL则通过--remote sc://host:port/;...传入CLI 解析器 SparkConnectClientParser.scala 的 usage 信息完整列出了这些选项与连接字符串参数表一一对应。JDBCjdbc:spark:// 前缀Spark Connect 的 JDBC 驱动sql/connect/client/jdbc在建立Connection时会将jdbc:前缀剥除后把剩余部分交给SparkConnectClient.builder().connectionString(...)解析并注入userAgent(Spark Connect JDBC)。也就是说 JDBC URL 形如jdbc:spark://host:port/;use_ssltrue;token...底层依然遵循本文的sc://连接字符串语法。连接字符串的解析与通道构建流程将 Python 客户端的实现串联起来可以清楚看到从字符串到 gRPC Channel 的完整链路core.pyScheme 校验URL 必须以sc://开头否则抛出INVALID_CONNECT_URL。URI 标准化把sc://改写为http://交给urllib.parse.urlparse解析出 host、port、path、params 等组件。路径校验path 必须为空或/否则报错。参数提取从 params 段按;拆分每个条目必须形如keyvalue值经 URL 解码后存入参数表。默认值填充port 缺省取 15002其余参数在读取时通过getDefault返回各自默认值keepalive 四项、消息大小等。通道构建根据secureuse_ssl 或 token选择 insecure / local credentials / SSL credentials 三种分支将 keepalive 参数注入 gRPC 通道选项最终返回grpc.Channel。同一套连接字符串语义在 JVM 侧由 SparkConnectClient.scala 实现两份实现保持参数名与默认值一致这正是跨语言统一连接表面设计目标的具体体现。参考与延伸连接字符串规范文档sql/connect/docs/client-connection-string.mdPython 客户端实现python/pyspark/sql/connect/client/core.pyPython 客户端单元测试python/pyspark/sql/tests/connect/client/test_client.pyJVM CLI 解析器sql/connect/common/src/main/scala/org/apache/spark/sql/connect/client/SparkConnectClientParser.scalaSpark Connect 服务端 keepalive 集成测试sql/connect/server/src/test/scala/org/apache/spark/sql/connect/service/SparkConnectServiceKeepAliveSuite.scalaSpark Connect JDBC 连接实现sql/connect/client/jdbc/src/main/scala/org/apache/spark/sql/connect/client/jdbc/SparkConnectConnection.scalaSpark Connect 总体介绍docs/spark-connect-overview.md赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐Hexo-theme-obsidian高级功能Live2D看板娘与自定义鼠标样式配置Hexo theme obsidian高级功能Live2D看板娘与自定义鼠标样式配置 Hexo theme obsidian是一款深色Hexo主题它响应式设大数据数据分析批处理流处理机器学习图计算Apache Spark Connect 概览解耦式客户端-服务端架构的完整实战指南Apache Spark Connect 概览解耦式客户端 服务端架构的完整实战指南 导读 本文以 Apache Spark 官方文档《Spark Conne大数据数据分析批处理流处理机器学习图计算Apache Airflow Spark Connect 连接配置指南sc:// 协议、认证参数与安全实践Apache Airflow Spark Connect 连接配置指南sc:// 协议、认证参数与安全实践 Apache Airflow 的 Apache S后端任务调度工作流自动化数据编排批处理数据工程流程编排上一篇3分钟上手秒传链接一个网页搞定百度网盘转存、生成与格式互转下一篇别再手动存说说了GetQzonehistory 一键备份创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表