PySpark 升级迁移指南:从 1.x 到 4.3 的行为变更、弃用与兼容性选项全解析
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
PySpark 每次大版本升级都伴随着 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.1)
Spark 4.2 是行为变更非常密集的一个版本,核心关键词是Arrow 全面默认化与Spark Connect 行为对齐。
3.1 依赖门槛与平台支持
- PyArrow 最低版本从 15.0.0 提升到 18.0.0;
- PyPy 不再被官方支持。官方建议在 CPython 上运行 PySpark;继续使用 PyPy 可能遇到未经验证的问题。
3.2 列式数据交换与 Arrow 优化默认开启
- PySpark 与 JVM 之间的列式数据交换默认使用 Apache Arrow:
spark.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.enabled(L5100-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_LEGACY=1。该环境变量在 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 可空整型扩展 dtype(
Int8/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(此前是全部缺失才抛)。建议先确认标签存在、先过滤出存在的标签,或传入errors="ignore"。
四、升级到 PySpark 4.1(自 4.0)
4.1 依赖门槛提升
- Python 3.9 支持被移除;
- PyArrow 最低版本从 11.0.0 提升到 15.0.0;
- Pandas 最低版本从 2.0.0 提升到 2.2.0。
4.2 Arrow UDF 行为修正与还原开关
DataFrame.__getitem__在 Spark Connect Python 客户端中不再急切校验列名,恢复方式同样是设置PYSPARK_VALIDATE_COLUMN_NAME_LEGACY=1。- Arrow 优化的 Python UDF 现在支持 UDT(用户自定义类型)输入/输出,不再回退到普通 UDF。要恢复旧的回退行为,设置
spark.sql.execution.pythonUDF.arrow.legacy.fallbackOnUDT=true。该内部配置定义于 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.enabled(L5596-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.binaryAsBytes=false(定义见 SQLConf.scala,默认true)。4.1.0 之前各场景下BinaryType对应的 Python 类型如下表:
| 场景 | BinaryType对应的 Python 类型 |
|---|---|
| 未启用 Arrow 优化的普通 UDF 与 UDTF | bytearray |
| DataFrame API(Spark Classic 与 Spark Connect) | bytearray |
| 数据源 | bytearray |
| 启用 Arrow 优化且带 pandas 转换的 UDF 与 UDTF | bytes |
4.4 pandas API on Spark 的 ANSI 模式
compute.ansi_mode_support默认True时,pandas API on Spark 可在 ANSI 模式下工作;原有的保护开关compute.fail_on_ansi_mode仍然保留,但只在compute.ansi_mode_support=False时才生效。
五、升级到 PySpark 4.0(自 3.5)
5.1 依赖门槛提升
- Python 3.8 支持被移除;
- Pandas 最低版本从 1.0.5 提升到 2.0.0;
- NumPy 最低版本从 1.15 提升到 1.21;
- PyArrow 最低版本从 4.0.0 提升到 11.0.0。
5.2 pandas API on Spark 的 API 移除与改名清单
4.0 对 pandas API on Spark 做了一次大规模清理,主要分三类:
整体移除(改用替代 API)
| 移除项 | 替代方案 |
|---|---|
Int64Index、Float64Index | 直接使用Index |
DataFrame.iteritems/Series.iteritems | DataFrame.items/Series.items |
DataFrame.append/Series.append | ps.concat |
DataFrame.mad/Series.mad | 无 |
Index.factorize/Series.factorize的na_sentinel参数 | 改用use_na_sentinel |
DataFrame.koalas | DataFrame.pandas_on_spark |
DataFrame.to_koalas/DataFrame.to_pandas_on_spark | DataFrame.pandas_api |
pyspark.testing.assertPandasOnSparkEqual | pyspark.pandas.testing.assert_frame_equal |
DataFrame.to_spark_io | DataFrame.spark.to_spark_io |
Index.asi8 | Index.astype |
Index.is_type_compatible | Index.isin |
Index.is_monotonic/Series.is_monotonic | is_monotonic_increasing系列方法 |
DataFrame.get_dtype_counts | DataFrame.dtypes.value_counts() |
DataFrameGroupBy.backfill/.pad | DataFrameGroupBy.bfill/.ffill |
Index.is_all_dates | 无 |
DatatimeIndex.week/.weekofyear、Series.dt.week/.weekofyear | DatetimeIndex.isocalendar().week/Series.dt.isocalendar().week |
参数移除
DataFrame.between_time/Series.between_time:移除include_start、include_end,改用inclusive;Series.between:inclusive不再接受布尔值,改用"both"/"neither";DataFrame.plot/Series.plot:移除sort_columns;ps.read_csv/ps.read_excel:移除squeeze;DataFrame.info:移除null_counts,改用show_counts;DataFrame.to_latex/Series.to_latex:移除col_space;DataFrame.to_excel/Series.to_excel:移除encoding、verbose;read_csv/read_excel:移除mangle_dupe_cols;read_excel另移除convert_float;Categorical.*与CategoricalIndex.*系列方法:移除inplace参数;ps.date_range:移除closed参数。
行为变化
DatetimeIndex的day、month、year等日期属性由int64变为int32;Series.str.replace的regex参数默认值由True改为False;且当regex=True时,单个字符的pat被当作正则表达式而非字符串字面量;value_counts的结果名称固定为'count'(normalize=True时为'proportion'),索引以原对象命名;MultiIndex.append不再保留索引名;DataFrameGroupBy.agg传入列表时遵守as_index=False;DataFrame.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.enabled=true(内部配置,定义见 SQLConf.scala)。这与 3.4 引入的数组列推断开关spark.sql.pyspark.legacy.inferArrayTypeFromFirstElement.enabled(L7646-L7653)是一对,共同决定SparkSession.createDataFrame的集合类型推断粒度。 compute.ops_on_diff_frames默认开启(支持跨不同 DataFrame 的运算);恢复旧行为需设为false。DataFrame.collect不再返回YearMonthIntervalType的底层整数;恢复旧行为需设置环境变量PYSPARK_YM_INTERVAL_LEGACY=1。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.enabled=false禁用 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.enabled=true(内部配置,默认false,见 SQLConf.scala)。 Groupby.apply的func未指定返回类型且compute.shortcut_limit=0时,采样行数固定为 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.dtype;GroupBy.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_PATCH=1恢复旧行为。
6.2 自 PySpark 3.2 升级到 3.3
pyspark.pandas.sql方法遵循 Python 标准字符串格式化语法;恢复旧行为需设置PYSPARK_PANDAS_SQL_LEGACY=1。- pandas API on Spark 的
drop方法支持按 index 删除行,且默认改为按行删除而非按列。 - Pandas 最低版本从 0.23.2 提升到 1.0.5。
- SQL 数据类型的
repr返回值改为可通过eval还原出等价对象。
6.3 自 PySpark 3.1 升级到 3.2
sql、ml、spark_on_pandas模块的方法在参数类型不匹配时抛TypeError而非ValueError。- Python UDF、pandas UDF 与 pandas function API 的 traceback 默认简化,不再打印内部 Python worker 的堆栈;恢复旧行为需设置
spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled=false。 - 默认启用 pinned thread 模式:每个 Python 线程映射到对应的 JVM 线程,避免多个 Python 线程共享一个 JVM 线程的 thread-local。建议配合
pyspark.InheritableThread或pyspark.inheritable_thread_target使用,以正确继承 JVM 线程的可继承属性(如 local properties),并避免资源泄漏。恢复旧行为需设置PYSPARK_PIN_THREAD=false。
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.enabled=true等)要求 PyArrow ≥ 0.12.1。 SparkSession.builder.getOrCreate()不再尝试用 builder 中的配置更新已存在的SparkContext(与 Java/Scala API 自 2.3 起行为一致);如需更新配置,须在创建SparkSession之前完成。- Arrow 优化下,PyArrow > 0.11.0 时可通过
spark.sql.execution.pandas.convertToArrowArraySafely=true开启安全类型转换(默认false)。不同版本/配置下的行为:
| PyArrow 版本 | 整数溢出 | 浮点截断 |
|---|---|---|
| 0.11.0 及以下 | 抛错 | 静默允许 |
> 0.11.0,arrowSafeTypeConversion=false | 静默溢出 | 静默允许 |
> 0.11.0,arrowSafeTypeConversion=true | 抛错 | 抛错 |
createDataFrame(..., verifySchema=True)开始校验LongType(此前不校验,溢出时得到None);可通过verifySchema=False关闭校验。Row按命名参数构造时字段顺序与输入一致(Python ≥ 3.6),不再按字母序排序;需要恢复排序时设置环境变量PYSPARK_ROW_FIELD_SORTING_ENABLED=true(必须在所有 driver 与 executor 上保持一致,否则可能引发失败或错误结果)。pyspark.ml.param.shared.Has*混入类不再提供set*(self, value)方法,改用self.set(self.*, value)。
6.5 更早期版本(2.x 与 1.x)
- 2.3→2.4:Arrow 优化开启时,
toPandas与从 Pandas DataFrame 创建 DataFrame 默认允许回退到非优化路径(此前toPandas直接失败);可通过spark.sql.execution.arrow.fallback.enabled关闭回退。 - 2.3.0→2.3.1:Arrow 功能(
pandas_udf、toPandas/createDataFrame配合spark.sql.execution.arrow.enabled=true)被标记为实验性,不建议在生产使用。 - 2.2→2.3:pandas 相关功能要求 Pandas ≥ 0.19.2;时间戳行为改为尊重会话时区(恢复旧行为设
spark.sql.execution.pandas.respectSessionTimeZone=false);na.fill()/fillna支持布尔值替换 null;df.replace在to_replace不是字典时不允许省略value。 - 1.4→1.5:字符串解析列支持用点号(
.)限定列名或访问嵌套值(如df['table.column.nestedField']),但列名中含点号时必须用反引号转义(如table.column.with.dots.nested);withColumn支持新增或替换同名列。 - 1.0–1.2→1.3:Python 中使用 DataTypes 时必须构造实例(如
StringType()),不再引用单例。
七、迁移开关速查表
将全文涉及的恢复开关汇总如下,便于升级排查时快速定位(源码定义见 SQLConf.scala):
| 类型 | 键名 | 影响的版本变更 | 默认值 |
|---|---|---|---|
| SQL 配置 | spark.sql.execution.arrow.pyspark.enabled | 4.2 起默认启用 Arrow 列式交换 | true |
| SQL 配置 | spark.sql.execution.pythonUDF.arrow.enabled | 4.2 起 UDF 默认 Arrow 优化 | true |
| SQL 配置 | spark.sql.execution.pythonUDTF.arrow.enabled | 4.2 起 UDTF 默认 Arrow 优化 | true |
| SQL 配置 | spark.sql.execution.pythonUDF.arrow.legacy.fallbackOnUDT | 4.1 UDT 回退旧行为 | false |
| SQL 配置 | spark.sql.legacy.execution.pythonUDF.pandas.conversion.enabled | 4.1 恢复 pandas 转换 | false |
| SQL 配置 | spark.sql.legacy.execution.pythonUDTF.pandas.conversion.enabled | 4.1 恢复 pandas 转换 | false |
| SQL 配置 | spark.sql.execution.pyspark.binaryAsBytes | 4.1BinaryType→bytes映射 | true |
| SQL 配置 | spark.sql.execution.pandas.convertToArrowArraySafely | 4.1 安全类型转换 | true |
| SQL 配置 | spark.sql.pyspark.legacy.inferMapTypeFromFirstPair.enabled | 4.0 map schema 推断 | false |
| SQL 配置 | spark.sql.pyspark.legacy.inferArrayTypeFromFirstElement.enabled | 3.4 array schema 推断 | false |
| SQL 配置 | compute.ops_on_diff_frames | 4.0 跨 frame 运算 | true |
| SQL 配置 | spark.sql.ansi.enabled | 4.0 ANSI 模式 | true |
| SQL 配置 | spark.sql.execution.pyspark.udf.simplifiedTraceback.enabled | 3.2 简化 traceback | true |
| 环境变量 | PYSPARK_VALIDATE_COLUMN_NAME_LEGACY | 4.1/4.2 列名急切校验 | 未设置 |
| 环境变量 | PYSPARK_YM_INTERVAL_LEGACY | 4.0YearMonthIntervalType底层整数 | 未设置 |
| 环境变量 | PYSPARK_ENABLE_NAMEDTUPLE_PATCH | 3.4 namedtuple patch | 未设置 |
| 环境变量 | PYSPARK_PANDAS_SQL_LEGACY | 3.3pyspark.pandas.sql格式化 | 未设置 |
| 环境变量 | PYSPARK_PIN_THREAD | 3.2 pinned thread | 未设置 |
| 环境变量 | PYSPARK_ROW_FIELD_SORTING_ENABLED | 3.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.convertToArrowArraySafely=true下的溢出报错视为需要修复的数据问题而非可绕过的偶发异常。 - 用“恢复开关”做灰度:对存量作业,可按上表在提交参数(
spark-submit --conf或spark-defaults.conf)中临时恢复旧行为,逐项消除行为差异后再移除开关,避免“一步到位”式升级导致的问题难以定位。 - 重点回归 pandas API on Spark:4.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),仅供参考