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

