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
相关产品推荐
相关产品推荐

