Spark 3.3中cache()引发DataFrame行数及聚合结果异常的原因
Spark 3.3数据处理异常的可能原因分析
问题场景
在使用Spark 3.3时遇到两类异常:
- 函数内部返回的DataFrame执行
count()结果为1,但外部接收后执行count()结果变为2;缓存该DataFrame后结果恢复为1,且distinct().count()始终为1,同时幽灵行导致求和结果翻倍。 - 输入DataFrame均已缓存的情况下,
groupBy操作结果出现随机性:多次执行df7.count()结果在1和2之间波动,但上游df6.count()始终稳定为2。
核心代码片段
def create_df(df1, df2): df3 = df1.cache().select(...).join(df2.cache(), on=..., how='full') return df3 # count() 返回1 df4 = create_df(df1, df2) # count() 返回2
完整复现代码
from datetime import date from pyspark.sql import functions as f, SparkSession spark = SparkSession.builder.getOrCreate() # 完整脚本中三个DataFrame均已缓存,此处仅展示数据结构 df1 = spark.createDataFrame( [('1234', '123', '0001', '12345', '02', 1.0, None, None, None, 'EA')], schema='PON string, RN string, RP string, CPN string, ON string, RQ double, PQ double, PC double, CE string, U string' ) df2 = spark.createDataFrame( [('0000000000', '0000', '02', '12345', '1234', '42', 'EA', None, 1.0)], schema='RN string, RP string, IO string, CPN string, PON string, CM string, U string, AC double, AQ double' ) df5 = spark.createDataFrame( [('1234', '0010', '123456', date(2025, 1, 1), date(2024, 1, 1), None, False)], schema='PON string, P string, PPN string, ASD date, AFD date, RL integer, ISO boolean' ) # 预期返回2行 df3 = df1.join(df2, on=['RN', 'RP', 'CPN', 'PON', 'U'], how='full') df6 = df3.join(df5, on='PON', how='inner') # 合并两行数据 df7 = df6.select('PON', 'AC') \ .groupBy('PON') \ .agg(f.sum(f.col('AC')).alias('AC')) print(df7.count()) # 此处应输出1,但完整脚本中输出2
执行结果
>>> df7.count() 2 >>> df7.show(5, False) +----+----+ |PON |AC | +----+----+ |1234|null| +----+----+
随机性表现
>>> [df6.count() for _ in range(10)] [2, 2, 2, 2, 2, 2, 2, 2, 2, 2] >>> [df7.count() for _ in range(10)] [1, 2, 1, 2, 2, 2, 2, 1, 2, 2]
可能原因分析
1. 缓存与懒执行的冲突
Spark缓存为懒执行机制,仅当触发count()这类action操作时才会实际缓存数据。函数内调用df1.cache()、df2.cache()后执行count(),确实会缓存输入数据,但返回的df3执行计划仍依赖原始DataFrame的血缘关系。外部调用df4.count()时,执行计划重新解析可能导致join操作重复执行,而full join处理空值键时生成重复行;缓存df4后,执行计划被固化,重复计算被阻断,结果恢复正常。
2. Full Join的空值匹配异常
若join键存在空值,Spark full join会将所有空值视为相等匹配,但在特定分区或Shuffle场景下,空值的哈希分配可能出现异常,导致同一组空值行被多次输出产生幽灵行。这种情况会因Shuffle的随机性(如分区数、任务调度顺序)导致结果波动。
3. Spark 3.3版本已知Bug
Spark 3.3存在部分与join、缓存、聚合相关的已知问题:
- 缓存DataFrame参与full join时,执行计划可能错误重复读取缓存数据,导致结果重复
- GroupBy聚合处理null值时,部分场景会出现分区内数据重复统计,导致count结果不稳定
- Shuffle过程中空值数据的分区策略偶发异常,引发聚合结果波动
4. Shuffle的随机性影响
尽管输入DataFrame已缓存,但join和groupBy操作会触发Shuffle。Shuffle过程中数据分区分配存在随机性,处理含null值的行时,可能出现同一key的行被分配到多个分区,或聚合时重复统计同一行数据。而df6.count()稳定是因为仅统计总行数,不受聚合逻辑影响。
内容的提问来源于stack exchange,提问作者PHPirate
相关产品推荐
相关产品推荐

