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)得到的结果是符合预期的:
| id | date_start | date_end | status | date_end_of_last_close | rank |
|---|---|---|---|---|---|
| 1 | 2023-10-01 | 2023-10-02 | open | 2023-10-03 | 1 |
| 1 | 2023-10-02 | 2023-10-03 | close | 2023-10-03 | 2 |
| 2 | 2023-10-01 | 2023-10-02 | close | 2023-10-04 | 1 |
| 2 | 2023-10-02 | 2023-10-04 | close | 2023-10-04 | 2 |
| 3 | 2023-10-02 | 2023-10-04 | open | NULL | 1 |
| 3 | 2023-10-03 | 2023-10-06 | open | NULL | 2 |
过滤后的异常结果
当我对A执行过滤rank = 1并删除rank列后:
A_result = A.filter(f.col("rank") == 1).drop("rank") display(A_result)
得到的结果完全不对:
| id | date_start | date_end | status | date_end_of_last_close |
|---|---|---|---|---|
| 1 | 2023-10-01 | 2023-10-02 | open | NULL |
| 2 | 2023-10-01 | 2023-10-02 | close | 2023-10-02 |
| 3 | 2023-10-02 | 2023-10-04 | open | NULL |
对比测试:手动创建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)
这次得到的结果是完全符合预期的:
| id | date_start | date_end | status | date_end_of_last_close |
|---|---|---|---|---|
| 1 | 2023-10-01 | 2023-10-02 | open | 2023-10-03 |
| 2 | 2023-10-01 | 2023-10-02 | close | 2023-10-04 |
| 3 | 2023-10-02 | 2023-10-04 | open | NULL |
我的困惑与后续发现
现在我真的懵了,完全搞不懂为什么同一个逻辑,直接在A上过滤就出问题,手动重建B之后过滤就正常?我本来还要做一系列带窗口函数的变换和过滤操作,现在对PySpark的结果完全没信心了。
更新(2023-10-31):后来我发现这好像是PySpark版本的问题!在PySpark 3.4.1版本里运行这个脚本,结果是完全正常的,但在最新的3.5.0版本里就会出现上面的异常情况。
备注:内容来源于stack exchange,提问作者Dani
相关产品推荐
相关产品推荐

