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

如何用PySpark内置函数实现MapType列的聚合求和(无UDF)

用PySpark内置函数实现Map类型列的分组聚合求和

要实现按id分组后对Map中相同key的value求和,完全用PySpark内置函数就能搞定,不用UDF或applyInPandas,核心思路是把Map类型列拆分为键值对行,分组求和后再重新组合成Map。

完整代码示例

1. 构造测试数据

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, IntegerType, MapType
import pyspark.sql.functions as F

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

# 测试数据:id + Map类型的counter列
data = [
    (1, {"a": 2, "b": 3}),
    (1, {"a": 4, "c": 1}),
    (2, {"b": 5}),
    (2, {"a": 1, "b": 2})
]

schema = StructType([
    StructField("id", IntegerType(), nullable=True),
    StructField("counter", MapType(IntegerType(), IntegerType()), nullable=True)
])

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

2. 执行聚合逻辑

# 1. 拆分Map为key-value行
# 2. 按id+key分组求和
# 3. 按id聚合,把key和sum值重新组合成Map
sum_counter_df = df \
    .select("id", F.explode("counter").alias("key", "value")) \
    .groupBy("id", "key") \
    .agg(F.sum("value").alias("sum_value")) \
    .groupBy("id") \
    .agg(
        F.map_from_arrays(
            F.collect_list("key"),
            F.collect_list("sum_value")
        ).alias("sum_counter")
    )

sum_counter_df.show(truncate=False)

输出结果

+---+----------------+
|id |sum_counter     |
+---+----------------+
|1  |{a -> 6, b -> 3, c -> 1}|
|2  |{b -> 7, a -> 1}|
+---+----------------+

适配不同类型

如果你的Map键值是字符串类型,只需要修改MapType的参数为MapType(StringType(), IntegerType()),聚合逻辑完全通用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 05:50:27