You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark 3.3.0使用ps.concat时未复用缓存DataFrame的问题求助

PySpark 3.3.0中ps.concat未使用缓存DataFrame的问题分析与解决方法

问题描述

将作业升级到PySpark 3.3.0后,使用pyspark.pandas的ps.concat([df1, df2])合并缓存的ps.DataFrame时出现异常:合并后的DataFrame未使用缓存数据,而是重新读取源数据,导致数据源认证问题。该问题在PySpark 3.2.3中不存在。

以下是可复现问题的极简代码:

import pyspark.pandas as ps
import pyspark
from pyspark.sql import SparkSession

import sys
import os
os.environ["PYSPARK_PYTHON"] = sys.executable

spark = SparkSession.builder.appName('bug-pyspark3.3').getOrCreate()

df1 = ps.DataFrame(data={'col1': [1, 2], 'col2': [3, 4]}, columns=['col1', 'col2'])
df2 = ps.DataFrame(data={'col3': [5, 6]}, columns=['col3'])
cached_df1 = df1.spark.cache()
cached_df2 = df2.spark.cache()

cached_df1.count()
cached_df2.count()

merged_df = ps.concat([cached_df1,cached_df2], ignore_index=True)
merged_df.head()
merged_df.spark.explain()

版本执行计划差异

PySpark 3.2.3执行计划(使用缓存)

输出中存在InMemoryTableScan,说明合并操作读取了缓存数据:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [(cast(_we0#1300 as bigint) - 1) AS __index_level_0__#1298L, col1#1291L, col2#1292L, col3#1293L]
   +- Window [row_number() windowspecdefinition(_w0#1299L ASC NULLS FIRST, specifiedwindowframe(RowFrame, unboundedpreceding$(), currentrow$())) AS _we0#1300], [_w0#1299L ASC NULLS FIRST]
      +- Sort [_w0#1299L ASC NULLS FIRST], false, 0
         +- Exchange SinglePartition, ENSURE_REQUIREMENTS, [plan_id=356]
            +- Project [col1#1291L, col2#1292L, col3#1293L, monotonically_increasing_id() AS _w0#1299L]
               +- Union
                  :- Project [col1#941L AS col1#1291L, col2#942L AS col2#1292L, null AS col3#1293L]
                  :  +- InMemoryTableScan [col1#941L, col2#942L]
                  :        +- InMemoryRelation [__index_level_0__#940L, col1#941L, col2#942L, __natural_order__#946L], StorageLevel(disk, memory, deserialized, 1 replicas)
                  :              +- *(1) Project [__index_level_0__#940L, col1#941L, col2#942L, monotonically_increasing_id() AS __natural_order__#946L]
                  :                 +- *(1) Scan ExistingRDD[__index_level_0__#940L,col1#941L,col2#942L]
                  +- Project [null AS col1#1403L, null AS col2#1404L, col3#952L]
                     +- InMemoryTableScan [col3#952L]
                           +- InMemoryRelation [__index_level_0__#951L, col3#952L, __natural_order__#955L], StorageLevel(disk, memory, deserialized, 1 replicas)
                                 +- *(1) Project [__index_level_0__#951L, col3#952L, monotonically_increasing_id() AS __natural_order__#955L]
                                    +- *(1) Scan ExistingRDD[__index_level_0__#951L,col3#952L]

PySpark 3.3.0执行计划(未使用缓存)

输出中Union操作直接扫描源ExistingRDD,无InMemoryTableScan:

== Physical Plan ==
AttachDistributedSequence[__index_level_0__#771L, col1#762L, col2#763L, col3#764L] Index: __index_level_0__#771L
+- Union
   :- *(1) Project [col1#412L AS col1#762L, col2#413L AS col2#763L, null AS col3#764L]
   :  +- *(1) Scan ExistingRDD[__index_level_0__#411L,col1#412L,col2#413L]
   +- *(2) Project [null AS col1#804L, null AS col2#805L, col3#423L]
      +- *(2) Scan ExistingRDD[__index_level_0__#422L,col3#423L]

差异是否正常?

该版本间的差异不正常,属于PySpark 3.3.0中pyspark.pandas模块的回归问题。PySpark 3.3.0对ps.concat的实现逻辑做了调整,在处理带缓存的ps.DataFrame时,未正确保留缓存的依赖关系,而是直接引用了原始RDD,跳过了缓存的InMemoryRelation。

强制使用缓存的解决方法

方案1:转Spark DataFrame合并后转回ps.DataFrame

利用Spark原生的unionByName保留缓存依赖,操作如下:

# 将缓存后的ps.DataFrame转为Spark DataFrame
spark_df1 = cached_df1.to_spark()
spark_df2 = cached_df2.to_spark()

# 使用unionByName合并(允许缺失列)
merged_spark_df = spark_df1.unionByName(spark_df2, allowMissingColumns=True)

# 转回ps.DataFrame并重置索引
merged_df = ps.DataFrame(merged_spark_df).reset_index(drop=True)

merged_df.head()
merged_df.spark.explain() # 执行计划会显示使用缓存

方案2:合并后手动缓存结果

如果业务允许,直接对合并后的DataFrame执行缓存,后续操作将使用缓存数据:

merged_df = ps.concat([cached_df1,cached_df2], ignore_index=True).spark.cache()
merged_df.count() # 触发缓存写入
merged_df.head()
merged_df.spark.explain()

方案3:临时降级到PySpark 3.2.3

若业务对版本要求不高,暂时降级到PySpark 3.2.3可恢复原有缓存行为,直到官方修复该回归问题。


内容的提问来源于stack exchange,提问作者frco9

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 06:35:20