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

Spark中如何为更新后的PK与Timestamp添加对应JSON列

解决Spark中为每条最新主键记录添加完整JSON列的问题

你的错误出在lit(dfrank.toJSON.first())这行代码:dfrank.toJSON.first()只会取出转JSON后的第一行字符串,lit()会把这个字符串当作常量赋值给所有行,导致所有记录的JSON列内容完全一致,不是对应每条记录的完整数据。

正确实现方案(两种可选)

方案一:保留完整数据做窗口筛选(推荐)

调整你的步骤2和3,不要提前丢弃原数据的其他字段,直接基于完整DataFrame做窗口筛选,再生成JSON列:

  1. 复用Schema读取DataFrame(修正你原代码中Schema获取的错误写法):
// 先读取一次获取Schema,避免重复推断
val schema = spark.read.json("/HDFS path").schema
val df = spark.read.schema(schema).json("/HDFS path")
  1. 直接对完整数据做窗口筛选:
val dfLatest = df.withColumn("rank", row_number().over(
    Window.partitionBy($"PK1", $"PK2", $"PK3").orderBy($"ts".desc)
))
.filter($"rank" === 1)
.drop("rank") // 移除临时的rank列
  1. 添加完整JSON列:
    使用to_json(struct("*"))将当前行的所有字段转为JSON字符串:
val dfWithJson = dfLatest.withColumn("JSON", to_json(struct("*")))

方案二:关联原数据补全完整信息(适用于已保留筛选后小数据集的场景)

如果已经按原步骤得到了仅包含PK、ts的dfrank,可以通过关联原DataFrame补全完整数据,再生成JSON列:

// 关联原DataFrame,匹配对应PK和最新ts的完整记录
val dfFull = dfrank.join(df, Seq("PK1", "PK2", "PK3", "ts"), "inner")

// 生成完整JSON列
val dfWithJson = dfFull.withColumn("JSON", to_json(struct("*")))

关键说明

to_json(struct("*"))会针对每一行数据,将所有字段打包成结构体后转为JSON字符串,确保每行的JSON列都是当前记录的完整数据,而不是固定的某一行内容。

内容的提问来源于stack exchange,提问作者pavan hudekar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 11:40:51