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

Spark:利用Window partitionBy获取当前行前序行的首末经纬度值

没问题,我来帮你搞定这个需求!基于你已经处理好的DataFrame,我们可以通过Spark窗口函数的扩展用法,精准提取前序行中符合条件的首末经纬度值。

实现步骤详解

假设你的DataFrame名为processed_df,包含列:id, timestamp, count, precount, pretsp, lat, long。

1. 定义前序行的窗口范围

首先我们需要明确:要针对每个id分组,按timestamp排序,只取当前行之前的所有行作为前序数据集。用窗口函数定义这个范围:

import org.apache.spark.sql.expressions.Window

// 定义窗口:按id分组,按timestamp排序,范围是从开头到当前行的上一行
val pre_rows_window = Window.partitionBy("id")
  .orderBy("timestamp")
  .rowsBetween(Window.unboundedPreceding, Window.currentRow - 1)

如果你的pretsp已经明确了前序时间戳的固定范围(比如只取等于pretsp的行),可以把窗口的partitionBy改成partitionBy("id", "pretsp"),这样会更精准。

2. 筛选并提取首末经纬度

接下来我们要在窗口里筛选出和当前行timestamp相同,但经纬度不同的前序行,然后用first()和last()函数提取首末值。这里需要先把当前行的经纬度和timestamp存为单独列,方便窗口内的比较:

import org.apache.spark.sql.functions.{first, last, when, col, coalesce}

// 先把当前行的关键值存为独立列,避免窗口内混淆
val df_with_current = processed_df
  .withColumn("curr_timestamp", col("timestamp"))
  .withColumn("curr_lat", col("lat"))
  .withColumn("curr_long", col("long"))

// 计算目标列:符合条件的前序行首末经纬度
val result_df = df_with_current
  // 首个符合条件的纬度
  .withColumn("first_target_lat",
    first(
      when(
        (col("timestamp") === col("curr_timestamp")) && 
        ((col("lat") =!= col("curr_lat")) || (col("long") =!= col("curr_long"))),
        col("lat")
      ),
      ignoreNulls = true
    ).over(pre_rows_window)
  )
  // 首个符合条件的经度
  .withColumn("first_target_long",
    first(
      when(
        (col("timestamp") === col("curr_timestamp")) && 
        ((col("lat") =!= col("curr_lat")) || (col("long") =!= col("curr_long"))),
        col("long")
      ),
      ignoreNulls = true
    ).over(pre_rows_window)
  )
  // 末个符合条件的纬度
  .withColumn("last_target_lat",
    last(
      when(
        (col("timestamp") === col("curr_timestamp")) && 
        ((col("lat") =!= col("curr_lat")) || (col("long") =!= col("curr_long"))),
        col("lat")
      ),
      ignoreNulls = true
    ).over(pre_rows_window)
  )
  // 末个符合条件的经度
  .withColumn("last_target_long",
    last(
      when(
        (col("timestamp") === col("curr_timestamp")) && 
        ((col("lat") =!= col("curr_lat")) || (col("long") =!= col("curr_long"))),
        col("long")
      ),
      ignoreNulls = true
    ).over(pre_rows_window)
  )
  // 可选:如果没有符合条件的前序行,用当前行值填充
  .withColumn("first_target_lat", coalesce(col("first_target_lat"), col("curr_lat")))
  .withColumn("first_target_long", coalesce(col("first_target_long"), col("curr_long")))
  .withColumn("last_target_lat", coalesce(col("last_target_lat"), col("curr_lat")))
  .withColumn("last_target_long", coalesce(col("last_target_long"), col("curr_long")))
  // 清理临时列
  .drop("curr_timestamp", "curr_lat", "curr_long")

关键逻辑说明

  • when()函数负责过滤:只保留和当前行timestamp相同,但经纬度至少有一个不同的前序行
  • ignoreNulls = true确保聚合时跳过null值,避免影响结果
  • coalesce()是可选操作,用来处理没有符合条件前序行的场景,避免结果出现null

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:19:13