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

Spark时序数据:获取最后非空val对应的timestamp_ms问题

解决方案

要获取与last_lag_prev_val对应的原始时间戳,需用处理val的相同逻辑来处理timestamp_ms——通过last函数结合ignoreNulls=True,实现向前填充最近的非空时间戳。

修正后的代码

from pyspark.sql import Window, functions as F
from pyspark.sql import Row

# 示例数据
df = spark.createDataFrame([
    Row(timestamp_ms=1672531200000, val='19'),
    Row(timestamp_ms=1672532100000, val='20'),
    Row(timestamp_ms=1672533000000, val=None),
    Row(timestamp_ms=1672533900000, val='22'),
    Row(timestamp_ms=1672534800000, val=None),
    Row(timestamp_ms=1672535700000, val=None),
    Row(timestamp_ms=1672536600000, val='25'),
    Row(timestamp_ms=1672537500000, val='20'),
    Row(timestamp_ms=1672538400000, val='27')
])

# 定义窗口:按timestamp_ms排序,无分区
window_spec = Window.orderBy("timestamp_ms")

df_result = df.withColumn("lag_prev_val", F.lag("val").over(window_spec)) \
              .withColumn("last_lag_prev_val", F.last("lag_prev_val", ignoreNulls=True).over(window_spec)) \
              .withColumn("non_null_timestamp", F.when(F.col("val").isNotNull(), F.col("timestamp_ms"))) \
              .withColumn("last_lag_prev_time", F.lag(F.last("non_null_timestamp", ignoreNulls=True).over(window_spec)).over(window_spec)) \
              .drop("non_null_timestamp")

df_result.show()

逻辑说明

  1. non_null_timestamp列:仅在val非空时记录当前行的timestamp_ms,其余情况为null。
  2. last_lag_prev_time列:先通过last("non_null_timestamp", ignoreNulls=True)获取截至当前行最近的非空val对应的时间戳,再用lag偏移一行——因为last_lag_prev_val是上一个非空val,对应的时间戳也需匹配上一个非空val的时间。

输出结果

timestamp_msvallag_prev_vallast_lag_prev_vallast_lag_prev_time
167253120000019nullnullnull
16725321000002019191672531200000
1672533000000null20201672532100000
167253390000022null201672532100000
1672534800000null22221672533900000
1672535700000nullnull221672533900000
167253660000025null221672533900000
16725375000002025251672536600000
16725384000002720201672537500000

输出完全匹配理想结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 17:14:55