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()
逻辑说明
non_null_timestamp列:仅在val非空时记录当前行的timestamp_ms,其余情况为null。last_lag_prev_time列:先通过last("non_null_timestamp", ignoreNulls=True)获取截至当前行最近的非空val对应的时间戳,再用lag偏移一行——因为last_lag_prev_val是上一个非空val,对应的时间戳也需匹配上一个非空val的时间。
输出结果
| timestamp_ms | val | lag_prev_val | last_lag_prev_val | last_lag_prev_time |
|---|---|---|---|---|
| 1672531200000 | 19 | null | null | null |
| 1672532100000 | 20 | 19 | 19 | 1672531200000 |
| 1672533000000 | null | 20 | 20 | 1672532100000 |
| 1672533900000 | 22 | null | 20 | 1672532100000 |
| 1672534800000 | null | 22 | 22 | 1672533900000 |
| 1672535700000 | null | null | 22 | 1672533900000 |
| 1672536600000 | 25 | null | 22 | 1672533900000 |
| 1672537500000 | 20 | 25 | 25 | 1672536600000 |
| 1672538400000 | 27 | 20 | 20 | 1672537500000 |
输出完全匹配理想结果。
内容的提问来源于stack exchange,提问作者smurphy
相关产品推荐
相关产品推荐

