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
相关产品推荐
相关产品推荐

