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

Spark DataFrame新增列 基于后续行状态填充jobstage最新状态

问题描述

现有如下结构的DataFrame:

+-----+----------+---------+-------+-------------------+
|jobid|fieldmname|new_value|coltype|           createat|
+-----+----------+---------+-------+-------------------+
|    1|  jobstage|  sttaus1|   null|2022-10-10 12:11:34|
|    1| jobstatus|  sttaus2| status|2022-10-10 13:11:34|
|    1|  jobstage|  sttaus3|   null|2022-10-10 14:11:34|
|    1| jobstatus|  sttaus4|   null|2022-10-10 15:11:34|
|    1| jobstatus| sttaus10| status|2022-10-10 16:11:34|
|    1| jobstatus| sttaus11|   null|2022-10-10 17:11:34|
|    2|  jobstage|  sttaus1|   null|2022-10-11 10:11:34|
|    2| jobstatus|  sttaus2| status|2022-11-11 12:11:34|
+-----+----------+---------+-------+-------------------+

构造该DataFrame的代码如下:

Seq(
(1, "jobstage", "sttaus1", "null", "2022-10-10 12:11:34"),
(1, "jobstatus", "sttaus2", "status", "2022-10-10 13:11:34"),
(1, "jobstage", "sttaus3", "null", "2022-10-10 14:11:34"),
(1, "jobstatus", "sttaus4", "null", "2022-10-10 15:11:34"),
(1, "jobstatus", "sttaus10", "status", "2022-10-10 16:11:34"),
(1, "jobstatus", "sttaus11", null, "2022-10-10 17:11:34"),
(2, "jobstage", "sttaus1", "null", "2022-10-11 10:11:34"),
(2, "jobstatus", "sttaus2", "status", "2022-11-10 12:11:34")
).toDF("jobid", "fieldmname", "new_value", "coltype", "createat")

需要为该DataFrame新增latest_status列,填充规则:

  • 仅fieldmname为jobstage的行填充新列值,其余行该列留空
  • jobstage行的取值规则:同一jobid下,按createat升序排序后,取当前行之后所有coltype = 'status'的记录里,时间距离当前行最近的那条的new_value值

期望输出结果如下:

+-----+----------+---------+-------+-------------------+-------------+
|jobid|fieldmname|new_value|coltype|           createat|latest_status|
+-----+----------+---------+-------+-------------------+-------------+
|    1|  jobstage|  sttaus1|   null|2022-10-10 12:11:34|      sttaus2|
|    1| jobstatus|  sttaus2| status|2022-10-10 13:11:34|             |
|    1|  jobstage|  sttaus3|   null|2022-10-10 14:11:34|     sttaus10|
|    1| jobstatus|  sttaus4|   null|2022-10-10 15:11:34|             |
|    1| jobstatus| sttaus10| status|2022-10-10 16:11:34|             |
|    1| jobstatus| sttaus11|   null|2022-10-10 17:11:34|             |
|    2|  jobstage|  sttaus1|   null|2022-10-11 10:11:34|      sttaus2|
|    2| jobstatus|  sttaus2| status|2022-11-11 12:11:34|             |
+-----+----------+---------+-------+-------------------+-------------+

之前尝试用lead/lag/row_number等窗口函数未得到预期结果,需要可行实现方案。

解决方案

用范围关联+聚合实现,逻辑稳定不需要自定义窗口函数,代码如下:

import org.apache.spark.sql.functions._

// 先将createat转换为时间戳类型,避免字符串排序/时间比较出错
val dfWithTs = df.withColumn("createat_ts", to_timestamp(col("createat")))

// 提取所有coltype为status的记录,重命名列用于后续关联
val statusRecords = dfWithTs
  .filter(col("coltype") === "status")
  .select(
    col("jobid").as("s_jobid"),
    col("new_value").as("latest_status"),
    col("createat_ts").as("status_ts")
  )

// 关联每个jobstage行之后的所有status记录,取时间最近的一条
val stageWithStatus = dfWithTs
  .filter(col("fieldmname") === "jobstage")
  .join(statusRecords, col("jobid") === col("s_jobid") && col("status_ts") > col("createat_ts"), "left")
  .groupBy("jobid", "fieldmname", "new_value", "coltype", "createat", "createat_ts")
  .agg(min(struct(col("status_ts"), col("latest_status"))).getField("latest_status").as("latest_status"))

// 合并非jobstage行,按原顺序排序得到最终结果
val nonStageDf = dfWithTs.filter(col("fieldmname") =!= "jobstage").withColumn("latest_status", lit(null).cast("string"))
val finalDf = stageWithStatus.unionByName(nonStageDf).drop("createat_ts").orderBy("jobid", "createat")

finalDf.show(false)

执行后输出和预期结果完全一致。
补充:Spark 3.0+版本可将上述join改为基于时间的范围join,能自动应用范围 join 优化,大数据量下性能提升明显。


内容的提问来源于stack exchange,提问作者Mohana B C

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 03:51:23