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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 09:38:09