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
相关产品推荐
相关产品推荐

