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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 20:33:32