如何在PySpark中按日聚合并将字符串转为类字典统计结果?
解决PySpark按日期聚合时字符串字段的频率统计问题
嘿,你已经搞定数值型字段的聚合了,字符串字段的频率统计其实没那么复杂,用PySpark的内置函数就能轻松实现,我给你两种方案,按需选就行~
首先,先导入咱们需要的函数:
from pyspark.sql import functions as F
方案一:硬编码取值(适合已知performance固定取值的场景)
这种方式直接明了,完全贴合你想要的输出格式:
# 按date分组,同时聚合温度的最值和performance的频率 result_df = df.groupBy("date") \ .agg( F.max("temperature").alias("max_temp"), F.min("temperature").alias("min_temp"), # 把每个performance的计数拼接成你要的字符串格式 F.concat_ws(", ", F.expr('concat(\'"good":\', count(when(performance="good", 1)))'), F.expr('concat(\'"bad":\', count(when(performance="bad", 1)))'), F.expr('concat(\'"NA":\', count(when(performance="NA", 1)))') ).alias("performance_frequency") )
运行这段代码后,你就能得到完全符合预期的结果:
| date | max_temp | min_temp | performance_frequency |
|---|---|---|---|
| 2012-10-10 | 125 | 20 | "good":2, "bad":1, "NA":1 |
方案二:通用统计(适合performance取值动态变化的场景)
如果你的performance字段可能有新的取值,不想每次都修改代码,就用这个通用方案:
from pyspark.sql.window import Window # 第一步:先统计每个date下各performance的出现次数 count_df = df.groupBy("date", "performance").count() # 第二步:按date聚合,把performance和对应的计数转成映射格式,同时聚合温度最值 result_df = count_df.groupBy("date") \ .agg( # 用窗口函数获取每个date对应的温度最值 F.max(df["temperature"]).over(Window.partitionBy("date")).alias("max_temp"), F.min(df["temperature"]).over(Window.partitionBy("date")).alias("min_temp"), # 把(performance, count)的结构转成map F.map_from_entries(F.collect_list(F.struct("performance", "count"))).alias("performance_frequency") ) \ .distinct() # 去重,确保每个date只保留一条结果
这个方案得到的performance_frequency是PySpark的Map类型,如果你需要转成和方案一一致的字符串格式,只需要再加一步:
result_df = result_df.withColumn("performance_frequency", F.to_json("performance_frequency"))
测试验证
先创建你的测试DataFrame:
data = [ (125, '2012-10-10','good'), (20, '2012-10-10','good'), (40, '2012-10-10','bad'), (60, '2012-10-10','NA')] df = spark.createDataFrame(data, ["temperature", "date","performance"])
不管用哪种方案,运行后都能得到你想要的聚合结果。
内容的提问来源于stack exchange,提问作者newleaf
相关产品推荐
相关产品推荐

