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

PySpark中使用Window Functions与过滤操作时DataFrame结果不一致的问题求助

PySpark中使用Window Functions与过滤操作时DataFrame结果不一致的问题求助

大家好,我最近在PySpark里处理DataFrame时遇到了一个非常诡异的问题:当我使用带过滤条件的窗口函数做数据变换后,再进行过滤操作,得到的结果完全不符合预期。我整理了一个最小复现的例子,希望大家能帮我看看到底哪里出问题了!

复现代码与初始数据

from pyspark.sql import SparkSession
import pyspark.sql.functions as f
from pyspark.sql.window import Window as w
from datetime import datetime, date
from chispa.dataframe_comparer import assert_df_equality

# 初始化SparkSession
spark = SparkSession.builder.config("spark.sql.repl.eagerEval.enabled", True).getOrCreate()

# 创建基础DataFrame
df = spark.createDataFrame(
    [
        (1, date(2023, 10, 1), date(2023, 10, 2), "open"),
        (1, date(2023, 10, 2), date(2023, 10, 3), "close"),
        (2, date(2023, 10, 1), date(2023, 10, 2), "close"),
        (2, date(2023, 10, 2), date(2023, 10, 4), "close"),
        (3, date(2023, 10, 2), date(2023, 10, 4), "open"),
        (3, date(2023, 10, 3), date(2023, 10, 6), "open"),
    ],
    schema="id integer, date_start date, date_end date, status string"
)

# 定义两个窗口分区
partition = w.partitionBy("id").orderBy("date_start", "date_end").rowsBetween(w.unboundedPreceding, w.unboundedFollowing)
partition2 = w.partitionBy("id").orderBy("date_start", "date_end")

# 创建DataFrame A:添加窗口计算列和排名列
A = df.withColumn(
    "date_end_of_last_close",
    f.max(f.when(f.col("status") == "close", f.col("date_end"))).over(partition)
).withColumn(
    "rank",
    f.row_number().over(partition2)
)

DataFrame A的正常结果

运行后display(A)得到的结果是符合预期的:

iddate_startdate_endstatusdate_end_of_last_closerank
12023-10-012023-10-02open2023-10-031
12023-10-022023-10-03close2023-10-032
22023-10-012023-10-02close2023-10-041
22023-10-022023-10-04close2023-10-042
32023-10-022023-10-04openNULL1
32023-10-032023-10-06openNULL2

过滤后的异常结果

当我对A执行过滤rank = 1并删除rank列后:

A_result = A.filter(f.col("rank") == 1).drop("rank")
display(A_result)

得到的结果完全不对:

iddate_startdate_endstatusdate_end_of_last_close
12023-10-012023-10-02openNULL
22023-10-012023-10-02close2023-10-02
32023-10-022023-10-04openNULL

对比测试:手动创建DataFrame后的正常结果

我尝试把A的数据手动提取出来,重新创建一个DataFrame B,再执行同样的过滤操作:

# 定义Schema
schema = "id INT, date_start DATE, date_end DATE, status STRING, date_end_of_last_close DATE, rank INT"

# 手动填入A的数据
data = [
    (1, date(2023, 10, 1), date(2023, 10, 2), "open", date(2023, 10, 3), 1),
    (1, date(2023, 10, 2), date(2023, 10, 3), "close", date(2023, 10, 3), 2),
    (2, date(2023, 10, 1), date(2023, 10, 2), "close", date(2023, 10, 4), 1),
    (2, date(2023, 10, 2), date(2023, 10, 4), "close", date(2023, 10, 4), 2),
    (3, date(2023, 10, 2), date(2023, 10, 4), "open", None, 1),
    (3, date(2023, 10, 3), date(2023, 10, 6), "open", None, 2),
]

# 创建DataFrame B
B = spark.createDataFrame(data, schema)

# 执行同样的过滤
B_result = B.filter(f.col("rank") == 1).drop("rank")
display(B_result)

这次得到的结果是完全符合预期的:

iddate_startdate_endstatusdate_end_of_last_close
12023-10-012023-10-02open2023-10-03
22023-10-012023-10-02close2023-10-04
32023-10-022023-10-04openNULL

我的困惑与后续发现

现在我真的懵了,完全搞不懂为什么同一个逻辑,直接在A上过滤就出问题,手动重建B之后过滤就正常?我本来还要做一系列带窗口函数的变换和过滤操作,现在对PySpark的结果完全没信心了。

更新(2023-10-31):后来我发现这好像是PySpark版本的问题!在PySpark 3.4.1版本里运行这个脚本,结果是完全正常的,但在最新的3.5.0版本里就会出现上面的异常情况。


备注:内容来源于stack exchange,提问作者Dani

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.21 09:43:10