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

