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

Pyspark RDD按城市提取最大值(含并列)及多条件排序实现问题

Pyspark RDD 实现方案

核心问题说明

你之前的代码使用reduceByKey求最大值时,默认只会保留匹配到的第一条最大值记录,无法保留同一城市下多个日期的并列最大值条目,因此需要调整实现逻辑。

具体实现步骤

步骤1:计算每个城市的最大值

先单独聚合得到每个城市对应的最大值,方便后续和原记录做匹配:

# 若需按【城市+标识】维度取最大值,将下方key替换为 (x[0][0], x[0][2]) 即可
city_max_rdd = rdd.map(lambda x: (x[0][0], x[1])) \
                  .reduceByKey(max)

步骤2:关联原数据与城市最大值

将原RDD转换为城市为key的结构,和上一步得到的城市最大值RDD做关联:

# 若上一步用了城市+标识作为key,此处map的key也同步替换为 (x[0][0], x[0][2])
joined_rdd = rdd.map(lambda x: (x[0][0], x)) \
                .join(city_max_rdd)

步骤3:过滤出所有最大值记录

保留原记录数值等于对应城市最大值的条目,还原为原数据结构:

filtered_rdd = joined_rdd.filter(lambda x: x[1][0][1] == x[1][1]) \
                         .map(lambda x: x[1][0])

步骤4:按指定规则排序

直接通过多元组设置排序优先级,无需多次调用sortBy:

# 排序优先级对应:①数值降序 ②日期升序 ③城市名称降序
sorted_rdd = filtered_rdd.sortBy(
    lambda x: (x[1], x[0][1], x[0][0]),
    ascending=(False, True, False)
)

# 查看结果
print(sorted_rdd.collect())

输出示例

以上代码处理你提供的示例数据,最终输出结果如下:

[(('City1', '2020-03-27', 'X1'), 44), 
(('City1', '2020-03-28', 'X1'), 44), 
(('City5', '2020-03-25', 'X5'), 15), 
(('City3', '2020-03-28', 'X3'), 15), 
(('City2', '2020-03-26', 'X2'), 14), 
(('City4', '2020-03-27', 'X4'), 5)]

内容的提问来源于stack exchange,提问作者JohnDoe34

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 04:15:09