多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Spark PySpark:spark.conf 运行时配置接口(RuntimeConfig)完全解析

Apache Spark PySpark:spark.conf 运行时配置接口(RuntimeConfig)完全解析 大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载spark.conf是 PySpark 中在运行时读取、设置和重置 Spark 与 Spark SQL 配置项的统一入口。本文基于python/pyspark/sql/conf.py的RuntimeConfig源码实现及其配套测试完整讲解set、get、getAll、unset、isModifiable五个 API 的语义、边界行为与异常约定并对比经典客户端与 Spark Connect 两种链路下的实现差异帮助你在脚本中安全地做动态调参与配置审计。一、RuntimeConfig 是什么spark.conf 的底层实现在 PySpark 中运行时配置接口通过SparkSession的conf属性暴露。从 session.py 可以看到该属性的定义cached_property def conf(self) - RuntimeConfig: Runtime configuration interface for Spark. This is the interface through which the user can get and set all Spark and Hadoop configurations that are relevant to Spark SQL. When getting the value of a config, this defaults to the value set in the underlying :class:SparkContext, if any. return RuntimeConfig(self._jsparkSession.conf())三个关键事实它是缓存属性cached_property同一个 session 上spark.conf返回同一个RuntimeConfig实例读取优先级get时会回退到底层SparkContext中已设置的值即通过SparkConf启动参数设置的配置同样可读与 Hadoop 配置的联动conf.py 中类文档明确说明 “Options set here are automatically propagated to the Hadoop configuration during I/O”即通过它设置的配置项会在 I/O 阶段自动传播到 Hadoop 配置。RuntimeConfig本身是一个薄封装——经典客户端模式下它直接包装一个 JVM 端对象class RuntimeConfig: User-facing configuration API, accessible through SparkSession.conf. Options set here are automatically propagated to the Hadoop configuration during I/O. .. versionchanged:: 3.4.0 Supports Spark Connect. def __init__(self, jconf: JavaObject) - None: Create a new RuntimeConfig that wraps the underlying JVM object. self._jconf jconf官方 API 文档页 python/docs/source/reference/pyspark.sql/configuration.rst 正是以 autosummary 方式汇总了pyspark.sql.conf.RuntimeConfig并通过 pyspark.sql 索引页 的configuration条目进入文档导航树。二、核心 API 逐一解析RuntimeConfig提供 5 个公开成员完整定义见 conf.py方法 / 属性签名引入版本说明setset(key: str, value: Union[str, int, bool]) - None2.0.0设置一个运行时配置项getget(key: str, defaultNone 占位) - Optional[str]2.0.0读取配置值支持显式默认值getAll属性返回Dict[str, str]4.0.0返回当前 conf 中所有配置项unsetunset(key: str) - None2.0.0重置某个配置项isModifiableisModifiable(key: str) - bool2.4.0判断该配置是否允许在当前会话修改2.1 set只接受 str / int / booldef set(self, key: str, value: Union[str, int, bool]) - None: Sets the given Spark runtime configuration property. Examples -------- spark.conf.set(key1, value1) self._jconf.set(key, value)从 test_conf.py 的测试可以确认类型处理的实际行为spark.conf.set(foo, True)→ 读回trueFalse→falsespark.conf.set(foo, 1)→ 读回字符串1配置值最终以字符串存储传入None会抛出IllegalArgumentException传入Decimal(1)等其他类型会抛异常——即value 只接受 str、int、bool 三种类型。2.2 get内置默认值与显式默认值的优先级get的语义是RuntimeConfig中最容易误解的部分源码逻辑conf.pydef get(self, key: str, default: Union[Optional[str], _NoValueType] _NoValue) - Optional[str]: self._check_type(key, key) if default is _NoValue: return self._jconf.get(key) else: if default is not None: self._check_type(default, default) return self._jconf.get(key, default)文档字符串给出了四类行为约定# 1) 未显式设置的键 提供 default → 返回 default spark.conf.get(non-existent-key, my_default) my_default # 2) 已显式设置的键 → 返回其值 spark.conf.set(my_key, my_value) spark.conf.get(my_key) my_value # 3) 未显式设置但存在内置默认值的键不传 default→ 返回内置默认值 spark.conf.unset(spark.sql.sources.partitionOverwriteMode) spark.conf.get(spark.sql.sources.partitionOverwriteMode) STATIC # 4) 传了 default → 该 default 会覆盖内置默认值而不只是兜底 spark.conf.get(spark.sql.sources.partitionOverwriteMode, DYNAMIC) DYNAMIC特别注意第 3、4 条的区分不传default时若该 key 有内置默认值built-in default返回的是内置默认值一旦显式传入default它就完全取代内置默认值。测试代码test_conf.py对这条语义做了双重验证# 未显式设置时返回内置默认值 STATIC self.assertEqual(spark.conf.get(spark.sql.sources.partitionOverwriteMode), STATIC) # 显式传入 None 时返回 Nonedefault 覆盖了内置默认值 self.assertEqual(spark.conf.get(spark.sql.sources.partitionOverwriteMode, None), None)另一个边界传入defaultNone是合法的返回None但传入非字符串的 default如数字会触发类型检查见下文_check_type。2.3 getAll4.0.0 新增的全量快照property def getAll(self) - Dict[str, str]: Returns all properties set in this conf. .. versionadded:: 4.0.0 return dict(self._jconf.getAllAsJava())返回一个普通 Pythondict。test_get_all 验证了其一致性初始快照非空且不含临时键set(foo, bar)后再取getAll条目数恰好 1 且包含新键unset后恢复原状。适合用于配置审计、会话调试时 dump 当前全部配置。2.4 unset重置后 get 的异常约定def unset(self, key: str) - None: Resets the configuration property for the given key. self._jconf.unset(key)文档示例明确了unset之后再次get的异常路径 spark.conf.set(my_key, my_value) spark.conf.get(my_key) my_value spark.conf.unset(my_key) spark.conf.get(my_key) Traceback (most recent call last): ... pyspark...SparkNoSuchElementException: ... The SQL config my_key cannot be found...即对一个既没有显式值、也没有内置默认值的 key 调用get不传 default会抛出pyspark.errors.SparkNoSuchElementException。测试中对应的断言是self.assertRaisesRegex(Exception, not.set, lambda: spark.conf.get(not.set))。2.5 isModifiable区分可动态修改与静态配置def isModifiable(self, key: str) - bool: Indicates whether the configuration property with the given key is modifiable in the current session. .. versionadded:: 2.4.0 return self._jconf.isModifiable(key)它回答的问题是这个配置项能不能在会话运行时被set修改测试用例test_conf.py给出了两个对照样本self.assertTrue(spark.conf.isModifiable(spark.sql.execution.arrow.maxRecordsPerBatch)) self.assertFalse(spark.conf.isModifiable(spark.sql.warehouse.dir))也就是说spark.sql.execution.arrow.maxRecordsPerBatch是运行时可改的而spark.sql.warehouse.dir属于静态配置只能在创建 session 时指定。这个判定的服务端实现在 SQLConf.scaladef isModifiable(key: String): Boolean { containsConfigKey(key) !isStaticConfigKey(key) }即“该 key 是已知配置项且不是静态配置键”。经典模式下RuntimeConfig.isModifiable只是转发给SQLConf见 classic/RuntimeConfig.scala。实践建议在动态调参脚本里可以先用isModifiable探测再set避免对静态配置项误改。静态配置的典型限制可参考 sql-arrow-cache-format.md 中的说明——部分执行层配置属于静态配置spark.conf.set会在运行时会话中拒绝修改。2.6 内部类型检查_check_typedef _check_type(self, obj: Any, identifier: str) - None: Assert that an object is of type str. if not isinstance(obj, str): raise PySparkTypeError( errorClassNOT_EXPECTED_TYPE, messageParameters{ arg_name: identifier, expected_type: str, arg_type: type(obj).__name__, }, )key永远是强校验的spark.conf.get(123)会抛PySparkTypeError错误类为NOT_EXPECTED_TYPE参数expected_typestr、arg_namekey、arg_typeint测试见 test_conf.py。default只在非None时校验。三、Spark Connect 链路下的实现差异RuntimeConfig的文档标注.. versionchanged:: 3.4.0 Supports Spark Connect.——同一套spark.conf写法在 Spark Connect 客户端下走的是完全不同的链路connect/conf.py 中的RuntimeConf通过 protobuf 的ConfigRequest操作与服务端通信。connect/session.py 中SparkSession.conf返回的正是RuntimeConfproperty def conf(self) - RuntimeConf: return RuntimeConf(self.client)各操作到 proto 消息的映射关系Python APIproto 操作ConfigRequest.Operation返回处理set(key, value)ConfigRequest.Set(pairs[KeyValue])将服务端warnings逐个转成warnings.warnget(key)ConfigRequest.Get(keys[key])取result.pairs[0][1]get(key, default)ConfigRequest.GetWithDefault(pairs[KeyValue])同上getAllConfigRequest.GetAll()把result.pairs收集为dictunset(key)ConfigRequest.Unset(keys[key])输出服务端 warningsisModifiable(key)ConfigRequest.IsModifiable(keys[key])字符串true/false转 bool其他值抛PySparkValueError几个值得注意的实现细节均以 connect/conf.py 源码为准类型转换前置set与内部_set_all在客户端就把bool转成true/false、int转成字符串proto 里传输的永远是字符串文档继承RuntimeConf的每个方法都通过set.__doc__ PySparkRuntimeConfig.set.__doc__之类的方式直接复用经典RuntimeConfig的文档字符串保证两套前端文档一致静默批量设置_set_all(configs, silent)支持一次性Set多对KeyValuesilent标志由 proto 的ConfigRequest.Set.silent字段承载异常等价性测试 test_parity_python_worker_env.py 的注释指出写入端拒绝非法变量的行为统一收敛在RuntimeConfig.setSQLSET也会经过同一条路径——即经典与 Connect 两条前端在语义上保持 parity对等。四、可复制的完整使用示例综合上述语义下面是一个覆盖全部 API 的最小完整示例经典本地模式from pyspark.sql import SparkSession spark SparkSession.builder.appName(conf-demo).getOrCreate() # 1. 设置与读取 spark.conf.set(spark.sql.shuffle.partitions, 200) # int 会自动转字符串 spark.conf.set(my_custom_key, my_value) print(spark.conf.get(my_custom_key)) # - my_value # 2. 内置默认值 vs 显式默认值 print(spark.conf.get(spark.sql.sources.partitionOverwriteMode)) # 未显式设置 - 返回内置默认值 STATIC print(spark.conf.get(spark.sql.sources.partitionOverwriteMode, DYNAMIC)) # 显式传 default - 覆盖内置默认值返回 DYNAMIC print(spark.conf.get(spark.sql.sources.partitionOverwriteMode, None)) # defaultNone - 返回 None # 3. 未设置的键 print(spark.conf.get(not.set, fallback)) # - fallback # spark.conf.get(not.set) # 不传 default 会抛 SparkNoSuchElementException # 4. 是否允许运行时修改 print(spark.conf.isModifiable(spark.sql.shuffle.partitions)) # True print(spark.conf.isModifiable(spark.sql.warehouse.dir)) # False静态配置 # 5. 全量快照与重置 snapshot spark.conf.getAll # Dict[str, str] spark.conf.unset(my_custom_key) # spark.conf.get(my_custom_key) # 抛 SparkNoSuchElementException该模块自带的 doctest 入口 conf.py_test()也是以SparkSession.builder.master(local[4])启动后对整模块做doctest.testmod上述示例行为与之一致。五、关键文件索引内容路径文档页autosummary 入口python/docs/source/reference/pyspark.sql/configuration.rst经典客户端实现python/pyspark/sql/conf.pySparkSession.conf属性python/pyspark/sql/session.pySpark Connect 客户端实现python/pyspark/sql/connect/conf.py行为验证测试python/pyspark/sql/tests/test_conf.py服务端RuntimeConfigScalasql/core/src/main/scala/org/apache/spark/sql/classic/RuntimeConfig.scalaisModifiable判定逻辑sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala完整配置项清单SQLdocs/configuration.md适用前提与限制本文所有行为描述均以当前仓库源码为准覆盖版本标注来自文档字符串set/get/unset自 2.0.0、isModifiable自 2.4.0、getAll自 4.0.0、Spark Connect 支持自 3.4.0。经典与 Connect 两种前端在本文涉及的读写、默认值、异常与可修改性语义上保持一致但服务端的具体行为以实际部署的 Spark 版本为准。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐Nitro 配置与运行时配置完全指南从 nitro.config.ts 到 runtimeConfig 的实战解析Nitro 配置与运行时配置完全指南从 nitro.config.ts 到 runtimeConfig 的实战解析 本指南基于 Nitro v3 的配置体系AI 技能人工智能Apache Iceberg Spark 配置完全指南Catalogs、SQL Extensions 与运行时参数详解Apache Iceberg Spark 配置完全指南Catalogs、SQL Extensions 与运行时参数详解 导读 本文基于 Apache Iceb数据湖大数据数据存储Nuxt 运行时配置迁移指南从 process.env 全面迁移到 runtimeConfig useRuntimeConfigNuxt 运行时配置迁移指南从 process.env 全面迁移到 runtimeConfig useRuntimeConfig 导读 本文讲解如何将 N前端后端Web框架SSR创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表