ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

PySpark 升级迁移指南:从 1.x 到 4.3 的行为变更、弃用与兼容性选项全解析

PySpark 升级迁移指南:从 1.x 到 4.3 的行为变更、弃用与兼容性选项全解析 PySpark 升级迁移指南从 1.x 到 4.3 的行为变更、弃用与兼容性选项全解析【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/sparkPySpark 每次大版本升级都伴随着 Python 依赖版本门槛的提升、Arrow 列式优化的逐步默认化以及 pandas API on Spark 与 Spark Connect 客户端的一系列行为变更。本文以仓库内 pyspark_upgrade.rst 为骨架逐版本梳理从 4.2→4.3 直至 1.3 的官方迁移要点并结合 SQLConf.scala 等源码中的配置定义、错误条件定义与测试用例说明每个变更背后的实现细节、恢复旧行为的开关帮助你在升级前完成兼容性评估与回归测试。一、升级总览读懂这篇指南的方式PySpark 的升级文档按“从版本 X 升级到版本 Y”的粒度组织每一条变更通常属于以下四类之一运行时依赖门槛变化Python、Pandas、PyArrow、NumPy 的最低支持版本被提高甚至某些解释器如 PyPy不再被官方支持默认行为切换某个功能从默认关闭变为默认开启典型代表是 Arrow 列式数据交换与 Arrow 优化的 UDF/UDTF同时保留显式关闭的配置项API 移除或参数弃用pandas API on Spark 中大量对齐 pandas 新版本的 API 被移除或更名语义修正错误类型、类型映射、schema 推断、空值处理等行为向 pandas/Arrow 标准靠拢需要业务代码相应调整。绝大多数行为变更都提供了“恢复旧行为”的开关——要么是spark.*SQL 配置项要么是PYSPARK_*环境变量。这是升级时最重要的逃生通道。二、升级到 PySpark 4.3自 4.2本版本最直接的变更是运行时环境的收紧Python 3.10 支持被移除。PySpark 4.3 不再支持 Python 3.10请确保运行环境驱动与执行器两侧的 Python 版本一致使用受支持的 Python 版本3.11 及以上具体以当前版本发布说明为准。这是 4.2 移除 Python 3.9、4.0 移除 Python 3.8 之后PySpark 持续清理老旧 Python 版本策略的延续。三、升级到 PySpark 4.2自 4.1Spark 4.2 是行为变更非常密集的一个版本核心关键词是Arrow 全面默认化与Spark Connect 行为对齐。3.1 依赖门槛与平台支持PyArrow 最低版本从 15.0.0 提升到 18.0.0PyPy 不再被官方支持。官方建议在 CPython 上运行 PySpark继续使用 PyPy 可能遇到未经验证的问题。3.2 列式数据交换与 Arrow 优化默认开启PySpark 与 JVM 之间的列式数据交换默认使用 Apache Arrowspark.sql.execution.arrow.pyspark.enabled默认值变为true。要恢复旧的非 Arrow行式数据交换将该配置设为false。该配置在源码中定义于 SQLConf.scala其文档明确指出该优化作用于两处pyspark.sql.DataFrame.toPandas以及输入为 Pandas DataFrame 或 NumPyndarray时的SparkSession.createDataFrame同时标出不支持的数据类型为ArrayTypeofTimestampType。它还有一个 3.0 起已弃用的前身spark.sql.execution.arrow.enabled通过fallbackConf继承其默认值true见同文件 L5081-L5086。在启用 Arrow 时如需进一步降低内存占用可配合实验性配置spark.sql.execution.arrow.pyspark.selfDestruct.enabledL5100-L5109使用 Arrow 的 self-destruct 与 split-blocks 选项以 CPU 时间换取内存。普通 Python UDF 默认使用 Arrow 优化spark.sql.execution.pythonUDF.arrow.enabled默认值变为true。恢复旧行为需显式设为false。该配置定义于 SQLConf.scala自 3.4.0 引入说明中特别注明“仅当函数接收至少一个参数时该优化才能启用”。普通 Python UDTF 默认使用 Arrow 优化spark.sql.execution.pythonUDTF.arrow.enabled默认值变为true自 3.5.0 引入定义见 SQLConf.scala。3.3 Spark Connect 客户端行为对齐DataFrame.__getattr__不再急切校验列名。在 Spark Connect Python 客户端中访问不存在的列不再在取属性时立即抛错。如需恢复旧行为设置环境变量PYSPARK_VALIDATE_COLUMN_NAME_LEGACY1。该环境变量在 python/pyspark/sql/connect/dataframe.py 中被读取当PYSPARK_VALIDATE_COLUMN_NAME_LEGACY为1或列名以__开头时才执行原校验逻辑。DataFrame[Stream]Reader/Writer.option与.options过滤None值None现在被当作“未设置”处理不再像以前那样把 Javanull转发到 JVM从而与 Spark Connect Python 客户端SPARK-49263及OptionUtils._set_opts行为一致。实操要点想把选项设为默认值直接省略或传None想显式设为空字符串则必须传。pandas UDF 接收可空整型列时使用扩展 dtype当 pandas UDF 的输入批次包含空值时可空整型列会以 pandas 可空整型扩展 dtypeInt8/Int16/Int32/Int64交付而不再是float64。依赖float64输入假设的 UDF 代码需要更新。相关配置spark.sql.execution.pythonUDF.pandas.preferIntExtensionDtype默认false见 SQLConf.scala可影响整型在 Pandas UDF 执行时的 dtype 选择。3.4 数据源与流式读取的严格化校验Python Data Source 返回的 Arrow 数据与声明 schema 类型不匹配时报DATA_SOURCE_RETURN_SCHEMA_MISMATCH。此前只有列数与列名不匹配会报该错4.2 起列类型不匹配同样触发。该错误条件登记在 python/pyspark/errors/error-conditions.json并由 python/pyspark/sql/worker/plan_data_source_read.py 在 worker 侧校验触发对应测试见 python/pyspark/sql/tests/test_python_datasource.py。修复方式让数据源返回的数据类型与声明 schema 严格一致。SimpleDataSourceStreamReader.read()返回非空批次但 end offset 未推进时报SIMPLE_STREAM_READER_OFFSET_DID_NOT_ADVANCE不再无限重放同一批次并膨胀预取缓存。错误条件同样登记在 error-conditions.json测试用例见 python/pyspark/sql/tests/test_python_streaming_datasource.py。修复方式确保返回的 end offset 至少推进到超过最后一条记录。3.5 其他行为变更SparkSession.createDataFrame从 NumPyndarray创建时要求 PyArrow而非 pandas数据会直接转换为 Arrow 而不是先经过 pandas。使用前请安装 PyArrow如果之前禁用 Arrow 且依赖基于 NumPy dtype 的 schema 推断需要复核推断结果——现在遵循 Arrow 的类型映射。pandas API on Spark 的DataFrame.drop与Series.drop语义对齐 pandas只要指定的标签中有一个缺失就抛KeyError此前是全部缺失才抛。建议先确认标签存在、先过滤出存在的标签或传入errorsignore。四、升级到 PySpark 4.1自 4.04.1 依赖门槛提升Python 3.9 支持被移除PyArrow 最低版本从 11.0.0 提升到 15.0.0Pandas 最低版本从 2.0.0 提升到 2.2.0。4.2 Arrow UDF 行为修正与还原开关DataFrame.__getitem__在 Spark Connect Python 客户端中不再急切校验列名恢复方式同样是设置PYSPARK_VALIDATE_COLUMN_NAME_LEGACY1。Arrow 优化的 Python UDF 现在支持 UDT用户自定义类型输入/输出不再回退到普通 UDF。要恢复旧的回退行为设置spark.sql.execution.pythonUDF.arrow.legacy.fallbackOnUDTtrue。该内部配置定义于 SQLConf.scala默认false。移除不必要的 pandas 实例转换当spark.sql.execution.pythonUDF.arrow.enabled启用时JVM 与 Python worker 之间的反序列化不再做额外的 pandas 转换导致“产出 schema 与声明 schema 不一致”时的类型强制转换type coercion行为发生变化。恢复旧行为需开启内部配置spark.sql.legacy.execution.pythonUDF.pandas.conversion.enabled默认false见 SQLConf.scala。UDTF 有完全对应的spark.sql.legacy.execution.pythonUDTF.pandas.conversion.enabledL5596-L5605。spark.sql.execution.pandas.convertToArrowArraySafely默认开启。启用后PyArrow 会在整数溢出、浮点截断、精度损失等不安全转换时抛错影响 Arrow 启用的 UDF/pandas_udf 的返回序列化以及 PySpark DataFrame 的创建。恢复旧行为需设为false。该配置定义于 SQLConf.scala 附近。4.3 BinaryType 与 Python bytes 的统一映射BinaryType默认统一映射为 Pythonbytes。恢复旧行为需设置spark.sql.execution.pyspark.binaryAsBytesfalse定义见 SQLConf.scala默认true。4.1.0 之前各场景下BinaryType对应的 Python 类型如下表场景BinaryType对应的 Python 类型未启用 Arrow 优化的普通 UDF 与 UDTFbytearrayDataFrame APISpark Classic 与 Spark Connectbytearray数据源bytearray启用 Arrow 优化且带 pandas 转换的 UDF 与 UDTFbytes4.4 pandas API on Spark 的 ANSI 模式compute.ansi_mode_support默认True时pandas API on Spark 可在 ANSI 模式下工作原有的保护开关compute.fail_on_ansi_mode仍然保留但只在compute.ansi_mode_supportFalse时才生效。五、升级到 PySpark 4.0自 3.55.1 依赖门槛提升Python 3.8 支持被移除Pandas 最低版本从 1.0.5 提升到 2.0.0NumPy 最低版本从 1.15 提升到 1.21PyArrow 最低版本从 4.0.0 提升到 11.0.0。5.2 pandas API on Spark 的 API 移除与改名清单4.0 对 pandas API on Spark 做了一次大规模清理主要分三类整体移除改用替代 API移除项替代方案Int64Index、Float64Index直接使用IndexDataFrame.iteritems/Series.iteritemsDataFrame.items/Series.itemsDataFrame.append/Series.appendps.concatDataFrame.mad/Series.mad无Index.factorize/Series.factorize的na_sentinel参数改用use_na_sentinelDataFrame.koalasDataFrame.pandas_on_sparkDataFrame.to_koalas/DataFrame.to_pandas_on_sparkDataFrame.pandas_apipyspark.testing.assertPandasOnSparkEqualpyspark.pandas.testing.assert_frame_equalDataFrame.to_spark_ioDataFrame.spark.to_spark_ioIndex.asi8Index.astypeIndex.is_type_compatibleIndex.isinIndex.is_monotonic/Series.is_monotonicis_monotonic_increasing系列方法DataFrame.get_dtype_countsDataFrame.dtypes.value_counts()DataFrameGroupBy.backfill/.padDataFrameGroupBy.bfill/.ffillIndex.is_all_dates无DatatimeIndex.week/.weekofyear、Series.dt.week/.weekofyearDatetimeIndex.isocalendar().week/Series.dt.isocalendar().week参数移除DataFrame.between_time/Series.between_time移除include_start、include_end改用inclusiveSeries.betweeninclusive不再接受布尔值改用both/neitherDataFrame.plot/Series.plot移除sort_columnsps.read_csv/ps.read_excel移除squeezeDataFrame.info移除null_counts改用show_countsDataFrame.to_latex/Series.to_latex移除col_spaceDataFrame.to_excel/Series.to_excel移除encoding、verboseread_csv/read_excel移除mangle_dupe_colsread_excel另移除convert_floatCategorical.*与CategoricalIndex.*系列方法移除inplace参数ps.date_range移除closed参数。行为变化DatetimeIndex的day、month、year等日期属性由int64变为int32Series.str.replace的regex参数默认值由True改为False且当regexTrue时单个字符的pat被当作正则表达式而非字符串字面量value_counts的结果名称固定为countnormalizeTrue时为proportion索引以原对象命名MultiIndex.append不再保留索引名DataFrameGroupBy.agg传入列表时遵守as_indexFalseDataFrame.stack保证保留既有列顺序不再按字典序排序对 decimal 类型对象应用astype时缺失值的转换结果由False变为True日期时间别名Y、M、H、T、S被弃用改用YE、ME、h、min、s。5.3 schema 推断、通配导入与 ANSI 模式map 列的 schema 推断改为合并所有键值对的 schema。要恢复“仅从第一个非空键值对推断”的旧行为设置spark.sql.pyspark.legacy.inferMapTypeFromFirstPair.enabledtrue内部配置定义见 SQLConf.scala。这与 3.4 引入的数组列推断开关spark.sql.pyspark.legacy.inferArrayTypeFromFirstElement.enabledL7646-L7653是一对共同决定SparkSession.createDataFrame的集合类型推断粒度。compute.ops_on_diff_frames默认开启支持跨不同 DataFrame 的运算恢复旧行为需设为false。DataFrame.collect不再返回YearMonthIntervalType的底层整数恢复旧行为需设置环境变量PYSPARK_YM_INTERVAL_LEGACY1。from pyspark.sql.functions import *不再导入非函数对象。DataFrame、Column、StructType等需要分别从pyspark.sql、pyspark.sql.types等模块显式导入。ANSI 模式与 pandas API on Spark 的冲突处理Spark 4.0 默认启用 ANSI 模式而 pandas API on Spark 在 ANSI 模式下无法正常工作并会抛异常。两种处理方式显式设置spark.sql.ansi.enabledfalse禁用 ANSI 模式或将 pandas-on-spark 选项compute.fail_on_ansi_mode设为False强制运行可能引发意外行为。六、更早版本升级要点3.5 及以前6.1 自 PySpark 3.3 升级到 3.4数组列的 schema 推断改为合并所有元素的 schema恢复“仅从第一个元素推断”需设置spark.sql.pyspark.legacy.inferArrayTypeFromFirstElement.enabledtrue内部配置默认false见 SQLConf.scala。Groupby.apply的func未指定返回类型且compute.shortcut_limit0时采样行数固定为 2以保证 schema 推断准确。Index.insert越界时抛IndexError信息形如index {} is out of bounds for axis 0 with size {}对齐 pandas 1.4。Series.mode保留 series 名称Index.__setitem__先检查value是否为Column类型避免is_list_like抛意外的ValueError。astype(category)会根据原始数据 dtype 刷新categories.dtypeGroupBy.head/GroupBy.tail支持位置索引负数参数语义正确此前返回空 frame。groupby.apply的 schema 推断会先推断 pandas 类型以保证 dtype 精度Series.concat尊重sort参数DataFrame.__setitem__会复制并替换既有数组不再覆写原数组。SparkSession.sql与 pandas API on Spark 的sql新增args参数支持具名参数绑定到 SQL 字面量。pandas API on Spark 全面跟随 pandas 2.0相关 API 的弃用/移除详见 pandas 官方 release notes。移除对collections.namedtuple的自定义 monkey-patch默认使用cloudpickle若遇到相关 pickling 问题设置环境变量PYSPARK_ENABLE_NAMEDTUPLE_PATCH1恢复旧行为。6.2 自 PySpark 3.2 升级到 3.3pyspark.pandas.sql方法遵循 Python 标准字符串格式化语法恢复旧行为需设置PYSPARK_PANDAS_SQL_LEGACY1。pandas API on Spark 的drop方法支持按 index 删除行且默认改为按行删除而非按列。Pandas 最低版本从 0.23.2 提升到 1.0.5。SQL 数据类型的repr返回值改为可通过eval还原出等价对象。6.3 自 PySpark 3.1 升级到 3.2sql、ml、spark_on_pandas模块的方法在参数类型不匹配时抛TypeError而非ValueError。Python UDF、pandas UDF 与 pandas function API 的 traceback 默认简化不再打印内部 Python worker 的堆栈恢复旧行为需设置spark.sql.execution.pyspark.udf.simplifiedTraceback.enabledfalse。默认启用 pinned thread 模式每个 Python 线程映射到对应的 JVM 线程避免多个 Python 线程共享一个 JVM 线程的 thread-local。建议配合pyspark.InheritableThread或pyspark.inheritable_thread_target使用以正确继承 JVM 线程的可继承属性如 local properties并避免资源泄漏。恢复旧行为需设置PYSPARK_PIN_THREADfalse。6.4 自 PySpark 2.4 升级到 3.0使用 pandas 相关功能toPandas、从 pandas DataFrame 创建 DataFrame 等要求 Pandas ≥ 0.23.2使用 PyArrow 相关功能pandas_udf、toPandas、createDataFrame配合spark.sql.execution.arrow.enabledtrue等要求 PyArrow ≥ 0.12.1。SparkSession.builder.getOrCreate()不再尝试用 builder 中的配置更新已存在的SparkContext与 Java/Scala API 自 2.3 起行为一致如需更新配置须在创建SparkSession之前完成。Arrow 优化下PyArrow 0.11.0 时可通过spark.sql.execution.pandas.convertToArrowArraySafelytrue开启安全类型转换默认false。不同版本/配置下的行为PyArrow 版本整数溢出浮点截断0.11.0 及以下抛错静默允许 0.11.0arrowSafeTypeConversionfalse静默溢出静默允许 0.11.0arrowSafeTypeConversiontrue抛错抛错createDataFrame(..., verifySchemaTrue)开始校验LongType此前不校验溢出时得到None可通过verifySchemaFalse关闭校验。Row按命名参数构造时字段顺序与输入一致Python ≥ 3.6不再按字母序排序需要恢复排序时设置环境变量PYSPARK_ROW_FIELD_SORTING_ENABLEDtrue必须在所有 driver 与 executor 上保持一致否则可能引发失败或错误结果。pyspark.ml.param.shared.Has*混入类不再提供set*(self, value)方法改用self.set(self.*, value)。6.5 更早期版本2.x 与 1.x2.3→2.4Arrow 优化开启时toPandas与从 Pandas DataFrame 创建 DataFrame 默认允许回退到非优化路径此前toPandas直接失败可通过spark.sql.execution.arrow.fallback.enabled关闭回退。2.3.0→2.3.1Arrow 功能pandas_udf、toPandas/createDataFrame配合spark.sql.execution.arrow.enabledtrue被标记为实验性不建议在生产使用。2.2→2.3pandas 相关功能要求 Pandas ≥ 0.19.2时间戳行为改为尊重会话时区恢复旧行为设spark.sql.execution.pandas.respectSessionTimeZonefalsena.fill()/fillna支持布尔值替换 nulldf.replace在to_replace不是字典时不允许省略value。1.4→1.5字符串解析列支持用点号.限定列名或访问嵌套值如df[table.column.nestedField]但列名中含点号时必须用反引号转义如table.column.with.dots.nestedwithColumn支持新增或替换同名列。1.0–1.2→1.3Python 中使用 DataTypes 时必须构造实例如StringType()不再引用单例。七、迁移开关速查表将全文涉及的恢复开关汇总如下便于升级排查时快速定位源码定义见 SQLConf.scala类型键名影响的版本变更默认值SQL 配置spark.sql.execution.arrow.pyspark.enabled4.2 起默认启用 Arrow 列式交换trueSQL 配置spark.sql.execution.pythonUDF.arrow.enabled4.2 起 UDF 默认 Arrow 优化trueSQL 配置spark.sql.execution.pythonUDTF.arrow.enabled4.2 起 UDTF 默认 Arrow 优化trueSQL 配置spark.sql.execution.pythonUDF.arrow.legacy.fallbackOnUDT4.1 UDT 回退旧行为falseSQL 配置spark.sql.legacy.execution.pythonUDF.pandas.conversion.enabled4.1 恢复 pandas 转换falseSQL 配置spark.sql.legacy.execution.pythonUDTF.pandas.conversion.enabled4.1 恢复 pandas 转换falseSQL 配置spark.sql.execution.pyspark.binaryAsBytes4.1BinaryType→bytes映射trueSQL 配置spark.sql.execution.pandas.convertToArrowArraySafely4.1 安全类型转换trueSQL 配置spark.sql.pyspark.legacy.inferMapTypeFromFirstPair.enabled4.0 map schema 推断falseSQL 配置spark.sql.pyspark.legacy.inferArrayTypeFromFirstElement.enabled3.4 array schema 推断falseSQL 配置compute.ops_on_diff_frames4.0 跨 frame 运算trueSQL 配置spark.sql.ansi.enabled4.0 ANSI 模式trueSQL 配置spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled3.2 简化 tracebacktrue环境变量PYSPARK_VALIDATE_COLUMN_NAME_LEGACY4.1/4.2 列名急切校验未设置环境变量PYSPARK_YM_INTERVAL_LEGACY4.0YearMonthIntervalType底层整数未设置环境变量PYSPARK_ENABLE_NAMEDTUPLE_PATCH3.4 namedtuple patch未设置环境变量PYSPARK_PANDAS_SQL_LEGACY3.3pyspark.pandas.sql格式化未设置环境变量PYSPARK_PIN_THREAD3.2 pinned thread未设置环境变量PYSPARK_ROW_FIELD_SORTING_ENABLED3.0Row字段排序未设置八、升级实践建议先核对运行时版本确认 Python≥ 3.11自 4.3 起、PyArrow≥ 18.0.0自 4.2 起、Pandas≥ 2.2.0自 4.1 起、NumPy≥ 1.21自 4.0 起满足目标版本门槛并统一 driver 与 executor 的 Python 环境。将 Arrow 作为第一公民对待4.2 起 Arrow 已是 PySpark 数据交换与 UDF/UDTF 执行的默认路径。升级前先用小规模任务验证 Arrow 路径下的类型映射尤其是BinaryType、可空整型 dtype、decimal 与时间类型再把spark.sql.execution.pandas.convertToArrowArraySafelytrue下的溢出报错视为需要修复的数据问题而非可绕过的偶发异常。用“恢复开关”做灰度对存量作业可按上表在提交参数spark-submit --conf或spark-defaults.conf中临时恢复旧行为逐项消除行为差异后再移除开关避免“一步到位”式升级导致的问题难以定位。重点回归 pandas API on Spark4.0 的大规模 API 清理对依赖ps的代码影响最大建议在升级前用pylint/静态扫描找出已移除的 API 与参数并对照上文替换清单逐一改写。关注 Spark Connect 客户端差异__getattr__/__getitem__列名校验放宽与option(s)的None过滤意味着依赖“取属性即抛错”或“传None即写 null”的代码需要显式调整。数据源与流式应用确保 Python Data Source 返回的 Arrow 数据类型与声明 schema 一致、SimpleDataSourceStreamReader的 end offset 正确推进否则 4.2 起会以明确错误码DATA_SOURCE_RETURN_SCHEMA_MISMATCH、SIMPLE_STREAM_READER_OFFSET_DID_NOT_ADVANCE失败。这些错误码统一定义在 python/pyspark/errors/error-conditions.json便于在文档与日志中检索。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表