如何用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
相关产品推荐
相关产品推荐

