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

PySpark中按id聚合行并保留MapType列的实现方法

PySpark按ID聚合MapType列的解决方案

需求说明

现有包含id和map_values(Map类型)两列的PySpark DataFrame,需要按id分组,将同组的所有map合并为一个Map类型列,保留所有键值对(重复键会被后续map的覆盖)。

示例数据构造

先创建示例DataFrame:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, IntegerType, MapType, StringType

spark = SparkSession.builder.appName("MapAggregation").getOrCreate()

# 定义表结构
schema = StructType([
    StructField("id", IntegerType(), nullable=True),
    StructField("map_values", MapType(StringType(), StringType()), nullable=True)
])

# 示例数据
sample_data = [
    (1, {"a": "b", "c": "d"}),
    (1, {"e": "f", "g": "h"}),
    (2, {"a": "b"})
]

df = spark.createDataFrame(sample_data, schema=schema)
df.show(truncate=False)

原始数据输出:

+---+--------------------------+
|id |map_values                |
+---+--------------------------+
|1  |{a -> b, c -> d}          |
|1  |{e -> f, g -> h}          |
|2  |{a -> b}                  |
+---+--------------------------+

聚合实现方法

方法1:内置函数聚合(Spark 3.0+推荐)

利用collect_list收集同id的所有map,再通过aggregate结合map_concat依次合并,全程使用Spark内置函数,性能更优:

from pyspark.sql.functions import lit, aggregate, collect_list, map_concat

aggregated_df = df.groupBy("id").agg(
    aggregate(
        collect_list("map_values"),
        lit({}).cast(MapType(StringType(), StringType())),  # 初始化空map
        lambda acc, x: map_concat(acc, x)
    ).alias("map_values")
)

aggregated_df.show(truncate=False)

方法2:自定义UDF(兼容低版本Spark)

如果你的Spark版本低于3.0,可通过自定义UDF实现map合并:

from pyspark.sql.functions import udf

def merge_maps(map_list):
    merged_map = {}
    for m in map_list:
        merged_map.update(m)
    return merged_map

# 注册UDF
merge_udf = udf(merge_maps, MapType(StringType(), StringType()))

aggregated_df = df.groupBy("id").agg(
    merge_udf(collect_list("map_values")).alias("map_values")
)

aggregated_df.show(truncate=False)

最终输出结果

两种方法都会得到如下结果:

+---+----------------------------------------+
|id |map_values                              |
+---+----------------------------------------+
|1  |{a -> b, c -> d, e -> f, g -> h}        |
|2  |{a -> b}                                |
+---+----------------------------------------+

注意:若同id的map存在重复键,后续map的键值会覆盖之前的,这是map_concat和字典update的默认行为。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 18:01:18