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

PySpark按序列号提取最近180天数据异常:返回全量数据

按序列号提取最近180天数据的问题

问题概述

我需要从包含10年数据的数据集中,按每个serial_number提取最近180天的数据。该逻辑在测试数据集上运行正常,但在真实数据集里未得到预期结果。

测试数据验证过程

测试目标

从测试数据中按序列号提取最近2天的数据。

测试数据加载代码

import findspark
findspark.init()
import pyspark
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName('spark3.2show').getOrCreate()
print('Spark info :')
spark
url="https://gist.githubusercontent.com/JishanAhmed2019/8c6a6effa98a8fe5cd86dc59d5959a87/raw/8c192ab825ad8191517bc9c2425a723df745cc2d/RecentNBeforeFailure.csv"
from pyspark import SparkFiles
spark.sparkContext.addFile(url)
df=spark.read.csv(SparkFiles.get("RecentNBeforeFailure.csv"), header=True,sep='\t')

提取最近2天数据的实现代码

import pyspark.sql.functions as F
from pyspark.sql.window import Window
df.withColumn(
    'date', F.to_timestamp(F.col('date'), 'M/D/yyyy')
).withColumn("r", F.rank().over(Window.partitionBy("serial_number") \
                           .orderBy(F.col("date").desc()))) \
    .filter("r <=2") \
    .drop("r") \
    .show()

测试结果

执行上述代码后,成功提取到每个序列号对应的最近2天数据,结果符合预期。

真实数据集异常情况

在真实数据集上执行以下提取最近180天数据的代码时:

import pyspark.sql.functions as F
from pyspark.sql.window import Window
OldSpark180=OldSpark.withColumn(
    'date', F.to_timestamp(F.col('date'), 'M/D/yyyy')
).withColumn("r", F.rank().over(Window.partitionBy("serial_number") \
                           .orderBy(F.col("date").desc()))) \
    .filter("r <=180") \
    .drop("r") 

返回了全部10年的数据,未筛选出最近180天的记录。

Pandas中的简单实现方式

在Pandas中可以通过以下代码轻松实现按序列号提取最近N条数据的需求:

pandas_df.sort_values('date').groupby('serial_number').tail(2)

内容的提问来源于stack exchange,提问作者ForestGump

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 06:27:42