如何仅按flow列分组,查询2016年各贸易流向最热门商品?
解决Spark SQL查找各贸易流向最频繁商品的问题
你的查询存在两个核心问题:
- 外层
GROUP BY flow后,SELECT里的commodity未被聚合函数包裹,Spark SQL会直接报错——分组查询的返回字段只能是分组字段或聚合计算结果。 - 即便强行保留
commodity,MAX(quantity)只能拿到每个贸易流向的最大出现次数,无法对应到具体的商品。
正确的解法是用窗口函数给每个贸易流向内的商品按出现次数排序,再取排名第一的记录:
query = ''' WITH commodity_count AS ( SELECT flow, commodity, COUNT(*) AS quantity FROM transactions WHERE year = 2016 GROUP BY flow, commodity ), ranked_commodities AS ( SELECT flow, commodity, quantity, ROW_NUMBER() OVER (PARTITION BY flow ORDER BY quantity DESC) AS rn FROM commodity_count ) SELECT flow, commodity, quantity FROM ranked_commodities WHERE rn = 1 ''' spark.sql(query).show(10)
代码说明:
- 第一个CTE
commodity_count:先统计2016年每个「贸易流向+商品」组合的出现次数。 - 第二个CTE
ranked_commodities:用ROW_NUMBER()窗口函数,按flow分组,在每个组内按quantity降序给商品排名。 - 最后过滤出排名为1的记录,就是每个贸易流向里出现次数最多的商品。
如果遇到同一贸易流向有多个商品出现次数相同的情况,ROW_NUMBER()只会取其中一条;若想保留所有并列第一的商品,可把ROW_NUMBER()换成RANK()。
执行后得到的结果格式与你期望的一致,示例如下:
[('出口', ('甜饼干、华夫饼和威化饼', 24)), ('进口', ('烘焙预拌粉', 27)), ('再出口', ('其他包装大于1kg的合成胶水或粘合剂', 8)), ('再进口', ('鞋底橡胶/塑料、鞋面织物的非运动鞋', 5))]
内容的提问来源于stack exchange,提问作者user21970405
相关产品推荐
相关产品推荐

