Spark中如何为更新后的PK与Timestamp添加对应JSON列
解决Spark中为每条最新主键记录添加完整JSON列的问题
你的错误出在lit(dfrank.toJSON.first())这行代码:dfrank.toJSON.first()只会取出转JSON后的第一行字符串,lit()会把这个字符串当作常量赋值给所有行,导致所有记录的JSON列内容完全一致,不是对应每条记录的完整数据。
正确实现方案(两种可选)
方案一:保留完整数据做窗口筛选(推荐)
调整你的步骤2和3,不要提前丢弃原数据的其他字段,直接基于完整DataFrame做窗口筛选,再生成JSON列:
- 复用Schema读取DataFrame(修正你原代码中Schema获取的错误写法):
// 先读取一次获取Schema,避免重复推断 val schema = spark.read.json("/HDFS path").schema val df = spark.read.schema(schema).json("/HDFS path")
- 直接对完整数据做窗口筛选:
val dfLatest = df.withColumn("rank", row_number().over( Window.partitionBy($"PK1", $"PK2", $"PK3").orderBy($"ts".desc) )) .filter($"rank" === 1) .drop("rank") // 移除临时的rank列
- 添加完整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
相关产品推荐
相关产品推荐

