Spark结构化流中按含Map类型字段分组获取最新时间戳行
解决Spark中Map类型无法作为窗口分区键的问题
问题背景
Spark 3.2.2中无法直接将Map类型字段作为窗口分区、分组或连接的键,会触发报错grouping/join/window partition keys cannot be map type。要实现按brand、original_timestamp和features(Map类型)分组并获取每组最新arrival_timestamp的需求,需要转换Map的表示形式。
方案1:通用Map序列化方案(适配任意结构的Map)
将Map类型字段序列化为JSON字符串,用该字符串作为临时分区键,后续保留原Map字段即可。
import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window // 将Map类型的features转为JSON字符串,作为临时分组键 val dfWithStrKey = df.withColumn("features_str", to_json(col("features"))) // 定义窗口:按brand、original_timestamp、序列化字符串分区,按arrival_timestamp倒序排序 val windowSpec = Window.partitionBy("brand", "original_timestamp", "features_str") .orderBy(col("arrival_timestamp").desc) // 给每组行打排名,取排名为1的最新行,最后移除临时列 val result = dfWithStrKey .withColumn("row_rank", row_number().over(windowSpec)) .filter(col("row_rank") === 1) .drop("row_rank", "features_str") result.show(false)
方案2:固定键Map展开方案(性能更优,适用于Map键固定的场景)
如果features的键是固定的(比如示例中的f1、f2),可以将Map中的键值对展开为单独列,用这些列作为分区键。
import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window // 展开Map中的固定键为独立列 val dfWithFlattenedFeatures = df .withColumn("f1", col("features").getItem("f1")) .withColumn("f2", col("features").getItem("f2")) // 用展开后的字段作为分区键定义窗口 val windowSpec = Window.partitionBy("brand", "original_timestamp", "f1", "f2") .orderBy(col("arrival_timestamp").desc) // 过滤每组最新行,移除临时列 val result = dfWithFlattenedFeatures .withColumn("row_rank", row_number().over(windowSpec)) .filter(col("row_rank") === 1) .drop("row_rank", "f1", "f2") result.show(false)
说明
两种方案均支持Spark结构化流场景,用到的to_json、getItem、row_number及Window函数都兼容流处理模式。
内容的提问来源于stack exchange,提问作者Nab
相关产品推荐
相关产品推荐

