PySpark聚合问题:将分组统计结果转为单个JSON对象
问题:将分组统计结果转换为单个JSON对象
原始DataFrame:
df_s create_date city 0 1 1 1 2 2 2 1 1 3 1 4 4 2 1 5 3 2 6 4 3
需求:先按create_date和city分组统计数量,再为每个唯一的create_date生成以city为键、统计数为值的单个JSON对象(而非JSON数组)。当前代码执行后得到的是JSON数组,不符合预期,需调整实现逻辑。
解决方案
核心思路:用aggregate函数将分组后收集的多个map对象合并为单个map,再转换为JSON。
步骤1:分组统计数量
保留原逻辑,也可使用更简洁的写法:
from pyspark.sql import functions as f # 按create_date和city分组,统计每个组合的数量 df_grouped = df_s.groupBy("create_date", "city").agg(f.count("city").alias("count")) df_grouped.show()
输出:
+-----------+----+-----+ |create_date|city|count| +-----------+----+-----+ | 1| 4| 1| | 2| 1| 1| | 4| 3| 1| | 2| 2| 1| | 3| 2| 1| | 1| 1| 2| +-----------+----+-----+
步骤2:合并Map并转换为单个JSON对象
通过aggregate将每个create_date对应的多个map(city, count)合并成一个大的map,再用to_json转换为JSON字符串:
result_df = df_grouped.groupBy("create_date").agg( f.to_json( f.aggregate( f.collect_list(f.create_map(f.col("city"), f.col("count"))), f.lit({}).cast("map<string, int>"), lambda acc, x: f.map_concat(acc, x) ) ).alias("city_count_json") ) result_df.show(truncate=False)
输出:
+-----------+-----------------------+ |create_date|city_count_json | +-----------+-----------------------+ |3 |{"2":1} | |1 |{"4":1,"1":2} | |4 |{"3":1} | |2 |{"1":1,"2":1} | +-----------+-----------------------+
代码说明
collect_list(f.create_map(...)):收集每个create_date对应的所有(city, count)键值对组成的map列表aggregate(...):初始化一个空map,遍历列表中的每个map,用map_concat将其合并到累加器中,最终得到包含所有键值对的单个mapto_json(...):将合并后的map转换为JSON字符串
内容的提问来源于stack exchange,提问作者CezarySzulc
相关产品推荐
相关产品推荐

