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
相关产品推荐
相关产品推荐

