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

PySpark中如何基于近两年数据用ntile实现动态窗口排名?

问题:基于过去两年数据的PySpark NTile滚动排名

现有PySpark DataFrame包含id、date(季度日期)、score字段,当前代码实现了全局范围的score NTile排名。需求改为:针对每一条观测值,基于其日期过去两年内的所有数据进行NTile(10)分组排名。询问是否有优雅的实现方案,比如带条件的分区方式。

import pandas as pd
import numpy as np
from pyspark.sql import SparkSession

# Generate the data using pandas
num_observations = 1000
num_dates = 20

dates = pd.date_range(start='2017-01-01', periods=num_dates, freq='3M')
df_pandas = pd.DataFrame({
    'id': range(1, num_observations + 1),
    'date': sorted(list(dates) * (num_observations // num_dates + 1))[:num_observations],
    'score': np.random.normal(loc=400, scale=25, size=num_observations).astype(int)
})

# Create SparkSession
spark = SparkSession.builder.getOrCreate()

# Convert pandas DataFrame to PySpark DataFrame
df_spark = spark.createDataFrame(df_pandas)

from pyspark.sql import functions as F
from pyspark.sql.window import Window

window = Window.orderBy(F.col('score'))

# Global rank
df_spark = df_spark.withColumn('ranked_score', F.ntile(10).over(window))

# Window rank

解决方案:使用范围窗口(Range Window)实现时间滚动NTile

核心思路是利用PySpark的范围窗口,基于日期的时间范围(当前记录日期往前推2年)划定计算窗口,在窗口内对score排序后执行NTile分组。这种方案无需额外自连接,性能更优且逻辑清晰。

步骤1:统一日期类型并转换为时间戳

确保date列为标准日期类型,并转换为时间戳(long型),方便计算时间范围:

# 转换为日期类型(确保兼容)并生成时间戳列
df_spark = df_spark.withColumn("date", F.to_date(F.col("date"))) \
                   .withColumn("date_ts", F.col("date").cast("timestamp").cast("long"))

步骤2:定义滚动时间窗口

用INTERVAL计算2年对应的秒数(自动处理闰年等特殊情况),然后定义范围窗口:

# 计算2年对应的秒数
two_years_seconds = F.expr("EXTRACT(EPOCH FROM INTERVAL '2 YEARS')").cast("long")

# 定义滚动窗口:所有日期在当前记录日期前2年到当前日期的范围
# 若需按id单独计算自身的滚动排名,添加 partitionBy("id") 即可
rolling_window = Window.orderBy("date_ts") \
                       .rangeBetween(-two_years_seconds, 0)

步骤3:计算滚动NTile排名

在滚动窗口内按score排序,执行NTile(10)分组:

# 在滚动窗口内按score排序,生成10分位的排名
df_spark = df_spark.withColumn("rolling_ntile_10", 
                               F.ntile(10).over(rolling_window.orderBy("score")))

关键说明

  • 范围窗口 vs 行窗口:范围窗口基于时间数值筛选,不受数据缺失(如某季度无记录)影响,能精准匹配"过去两年"的需求;行窗口仅按行数筛选,无法适配时间范围逻辑。
  • 按id分区:若需求是每个id单独计算自己过去两年的NTile排名,只需在Window中添加partitionBy("id")即可。
  • 灵活性:修改INTERVAL '2 YEARS'可轻松调整时间范围(如'1 YEAR'、'3 YEARS')。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 20:05:18