如何对Spark DataFrame按ID分组聚合Time并保留全列?
如何在Spark分组聚合时保留非分组列?
嘿,我来帮你搞定这个问题!先理清楚你的场景和需求:
首先你有这样一个Spark DataFrame:
| Title | Status | Suite | ID | Time |
|---|---|---|---|---|
| KIM | Passed | ABC | 123 | 20 |
| KJT | Passed | ABC | 123 | 10 |
| ZXD | Passed | CDF | 123 | 15 |
| XCV | Passed | GHY | 113 | 36 |
| KJM | Passed | RTH | 456 | 45 |
| KIM | Passed | ABC | 115 | 47 |
| JY | Passed | JHJK | 8963 | 74 |
| KJH | Passed | SNMP | 256 | 47 |
| KJH | Passed | ABC | 123 | 78 |
| LOK | Passed | GHY | 456 | 96 |
| LIM | Passed | RTH | 113 | 78 |
| MKN | Passed | ABC | 115 | 74 |
| KJM | Passed | GHY | 8963 | 74 |
可以通过这段代码创建:
df = sqlCtx.createDataFrame( [ ('KIM', 'Passed', 'ABC', '123',20), ('KJT', 'Passed', 'ABC', '123',10), ('ZXD', 'Passed', 'CDF', '123',15), ('XCV', 'Passed', 'GHY', '113',36), ('KJM', 'Passed', 'RTH', '456',45), ('KIM', 'Passed', 'ABC', '115',47), ('JY', 'Passed', 'JHJK', '8963',74), ('KJH', 'Passed', 'SNMP', '256',47), ('KJH', 'Passed', 'ABC', '123',78), ('LOK', 'Passed', 'GHY', '456',96), ('LIM', 'Passed', 'RTH', '113',78), ('MKN', 'Passed', 'ABC', '115',74), ('KJM', 'Passed', 'GHY', '8963',74), ],('Title', 'Status', 'Suite', 'ID','Time') )
你的需求是按ID分组计算Time的均值,同时保留Title、Status、Suite列,但之前的代码只返回了ID和Time列——这是因为Spark分组后,只有分组列和你明确指定聚合操作的列会被保留,其他列如果没有对应的聚合逻辑,Spark不知道该如何处理组内的多个值,所以不会自动包含在结果里。
解决方案
根据你的期望输出,我们需要对Title、Status、Suite也添加聚合逻辑:
Status列所有值都是Passed,用任意聚合函数(比如first、max)都能得到正确结果;Title和Suite可以选择取组内第一个出现的值,或者取出现次数最多的众数(比如ID=123的Suite是ABC,出现次数最多)。
方案1:取组内第一个出现的非分组列
这个方案简单直接,完全匹配你的期望输出:
from pyspark.sql import functions as F result_df = df.groupBy("ID") \ .agg( F.first("Title").alias("Title"), F.first("Status").alias("Status"), F.first("Suite").alias("Suite"), F.mean("Time").alias("Time") ) \ .orderBy("ID") # 按ID排序,和你的期望输出顺序一致 result_df.show()
运行后得到的结果如下:
| Title | Status | Suite | ID | Time |
|---|---|---|---|---|
| KIM | Passed | ABC | 123 | 30.75 |
| XCV | Passed | GHY | 113 | 57.0 |
| KIM | Passed | ABC | 115 | 60.5 |
| KJH | Passed | SNMP | 256 | 47.0 |
| KJM | Passed | RTH | 456 | 70.5 |
| JY | Passed | JHJK | 8963 | 74.0 |
方案2:取Suite的众数(更贴合业务场景)
如果你的业务逻辑是每个ID对应最常见的Suite,可以用mode函数(Spark 3.0+支持):
result_df = df.groupBy("ID") \ .agg( F.first("Title").alias("Title"), F.first("Status").alias("Status"), F.mode("Suite").alias("Suite"), F.mean("Time").alias("Time") ) \ .orderBy("ID")
这样ID=123的Suite会正确取到出现次数最多的ABC,和你的期望一致。
内容的提问来源于stack exchange,提问作者Krishna
相关产品推荐
相关产品推荐

