You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何仅按flow列分组,查询2016年各贸易流向最热门商品?

解决Spark SQL查找各贸易流向最频繁商品的问题

你的查询存在两个核心问题:

  1. 外层GROUP BY flow后,SELECT里的commodity未被聚合函数包裹,Spark SQL会直接报错——分组查询的返回字段只能是分组字段或聚合计算结果。
  2. 即便强行保留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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 02:55:01