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

Spark Scala多列分组生成嵌套值映射的实现方法求助

解决Spark DataFrame嵌套Map聚合问题

嗨,我来帮你搞定这个需求!你之前遇到的Map只保留最后一个值的问题,是因为直接在顶层分组后构建Map时,相同year的条目会被覆盖。咱们得分两步来实现嵌套Map的聚合:

实现思路

  1. 先按三级分组收集value列表:把group_1、group_2和year作为分组键,收集每个年份对应的所有value到一个列表中,避免同一年份的数据被覆盖。
  2. 再按二级分组构建嵌套Map:把第一步得到的year和对应的value列表组合成键值对,最后聚合这些键值对成一个完整的Map。

Scala代码示例

import org.apache.spark.sql.functions._

// 第一步:按group_1、group_2、year分组,收集每个年份的value列表
val intermediateDF = df.groupBy("group_1", "group_2", "year")
  .agg(collect_list("value").alias("value_list"))

// 第二步:按group_1、group_2分组,将year和value_list转换为Map
val resultDF = intermediateDF.groupBy("group_1", "group_2")
  .agg(map_from_entries(collect_list(struct(col("year"), col("value_list")))).alias("values_map"))

// 查看结果
resultDF.show(truncate = false)

Python代码示例

from pyspark.sql import functions as F

# 第一步:按三级分组收集value列表
intermediate_df = df.groupBy("group_1", "group_2", "year") \
    .agg(F.collect_list("value").alias("value_list"))

# 第二步:按二级分组构建嵌套Map
result_df = intermediate_df.groupBy("group_1", "group_2") \
    .agg(F.map_from_entries(F.collect_list(F.struct(F.col("year"), F.col("value_list")))).alias("values_map"))

# 查看结果
result_df.show(truncate=False)

为什么之前的方法失效?

如果直接在groupBy("group_1", "group_2")后使用map(col("year"), collect_list(col("value"))),Spark会将每条数据的year和对应value单独处理,当同一个year有多条数据时,最终Map只会保留该year对应的最后一个条目。而通过先按year分组收集完整列表,再将year与列表作为整体构建Map,就能完美保留所有年份对应的全部value了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 21:52:51