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

如何为PySpark DataFrame添加窗口内最小值对应的时间戳列?

解决方案

针对你的需求,这里提供两种可行的实现方法,都能准确获取5分钟回溯窗口内Foo最小值对应的Timestamp:


方法一:利用窗口聚合+内置表达式提取

通过在窗口内收集Foo与Timestamp的配对数据,再筛选出最小值对应的时间戳,逻辑简洁直接。

代码实现

from pyspark.sql import functions as F

# 定义原有的5分钟回溯窗口
window_bw = Window.orderBy(F.col('Timestamp').cast('int')).rangeBetween(-5*60, 0)

# 1. 保留原逻辑添加min_value列
df = df.withColumn('min_value', F.min('Foo').over(window_bw))

# 2. 收集窗口内的(Foo, Timestamp)对,筛选出最小值对应的时间戳
df = df.withColumn('foo_ts_pairs', F.collect_list(F.struct('Foo', 'Timestamp')).over(window_bw)) \
       .withColumn('min_value_timestamp', 
                   F.expr("filter(foo_ts_pairs, x -> x.Foo = min_value)[0].Timestamp")) \
       .drop('foo_ts_pairs')

df.show(truncate=False)

说明

  • collect_list(F.struct('Foo', 'Timestamp'))会把窗口内每行的Foo和Timestamp打包成结构体,再收集为列表。
  • filter函数筛选出列表中Foo等于当前行min_value的结构体,取第一个元素的Timestamp(若窗口内有多个相同最小值,会取最早出现的时间戳)。

方法二:窗口排序标记最小值行+关联

通过在窗口内给行排序,标记出最小值对应的行,再通过关联把时间戳带回原表,适合需要灵活处理多最小值场景的情况。

代码实现

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

# 定义原有的5分钟回溯窗口
window_bw = Window.orderBy(F.col('Timestamp').cast('int')).rangeBetween(-5*60, 0)

# 1. 在窗口内按Foo升序、Timestamp升序给行编号,最小值行的编号为1
df_with_rank = df.withColumn('rank', 
                             F.row_number().over(window_bw.partitionBy().orderBy('Foo', 'Timestamp')))

# 2. 筛选出每个窗口的最小值行,保留时间戳和对应最小值
min_ts_df = df_with_rank.filter(F.col('rank') == 1) \
                        .select(F.col('Timestamp').alias('current_ts'), 
                                F.col('Timestamp').alias('min_value_timestamp'), 
                                F.col('Foo').alias('min_value'))

# 3. 关联回原表,匹配每行对应的窗口最小值时间戳
result_df = df.join(min_ts_df, on=F.col('Timestamp') == min_ts_df['current_ts'], how='left') \
              .drop('current_ts')

result_df.show(truncate=False)

说明

  • row_number()在窗口内按Foo从小到大排序,若Foo相同则按Timestamp排序,确保最小值对应的行编号为1。
  • 筛选出编号为1的行后,通过关联操作将最小值对应的时间戳映射到原表的每一行中。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 16:25:35