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

Spark按列最大值合并行生成Map类型的实现问题

你原有逻辑的窗口分区规则不符合需求,你需要先对每个id下的每个name取最新timestamp对应的value,再按id聚合为Map,完整实现如下:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.types._

val schema = StructType(List(
    StructField("id", LongType, nullable = true),
    StructField("name", StringType, nullable = true),
    StructField("value", StringType, nullable = true),
    StructField("timestamp", LongType, nullable = true)))

val myDF = spark.read.schema(schema).option("header", "true").option("delimiter", ",").csv("wasbs:///HdiSamples/HdiSamples/SensorSampleData/hvac/tru.csv")
val df = myDF.toDF("id","name","value","timestamp")

// 窗口按id、name分区,按timestamp倒序排序,取每个分区第一条就是对应name的最新记录
val windowSpec = Window.partitionBy("id", "name").orderBy(col("timestamp").desc)
val latestRecordDf = df.withColumn("rn", row_number().over(windowSpec))
  .where(col("rn") === 1)
  .drop("rn", "timestamp")

// 按id分组,将name、value收集为数组后转Map
val resultDf = latestRecordDf.groupBy("id")
  .agg(
    map_from_arrays(collect_list("name"), collect_list("value")).alias("event_properties")
  )

resultDf.show(false)

核心逻辑说明

  • 窗口分区调整为partitionBy("id", "name"):确保每个id下的每个name维度独立计算最新记录,避免不同name的timestamp互相干扰
  • row_number()排序后取第一条:比直接对比max(timestamp)更稳妥,避免同个id同个name出现多条相同最大timestamp的记录时重复
  • map_from_arrays函数:直接把收集到的name数组和value数组按位置映射为Map结构,符合输出要求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 08:06:04