【问题标题】:Pyspark display max value(S) and multiple sortingPyspark显示最大值(S)和多重排序
【发布时间】:2022-01-07 21:05:12
【问题描述】:

感谢您的帮助。使用 Pyspark(请不要使用 SQL)。所以我有一个存储为 RDD 对的元组列表:

[(('City1', '2020-03-27', 'X1'), 44),

(('City1', '2020-03-28', 'X1'), 44),

(('City3', '2020-03-28', 'X3'), 15),

(('City4', '2020-03-27', 'X4'), 5),

(('City4', '2020-03-26', 'X4'), 4),

(('City2', '2020-03-26', 'X2'), 14),

(('City2', '2020-03-25', 'X2'), 4),

(('City4', '2020-03-25', 'X4'), 1),

(('City1', '2020-03-29', 'X1'), 1),

(('City5', '2020-03-25', 'X5'), 15)]

以例如 ('City5', '2020-03-25', 'X5') 作为键,15 作为最后一对的值。

我想得到以下结果:

City1, X1, 2020-03-27, 44

City1, X1, 2020-03-28, 44

City5, X3, 2020-03-25, 15

City3, X3, 2020-03-28, 15

City2, X2, 2020-03-26, 14

City4, X4, 2020-03-27, 5

请注意结果显示:

  • 每个城市的最大值的键(这是最难的部分,如果同一城市在不同日期具有相似的最大值(值),则显示两次,我假设不能使用 ReduceByKey() 作为键是不是唯一的,可能是 GroupBy() 或 Filter() ?

  • 在以下排序/排序顺序中:

  1. 最大值递减
  2. 升序日期
  3. 降序城市名称(例如:City1)

所以我尝试了以下代码:

res = rdd2.map(lambda x: ((x[0][0],x[0][2]), (x[0][1], x[1])))
rdd3 = res.reduceByKey(lambda x1, x2: max(x1, x2, key=lambda x: x[1]))
rdd4 = rdd3.sortBy(lambda a: a[1][1], ascending=False)
rdd5 = rdd4.sortBy(lambda a: a[1][0])

虽然它确实给了我具有最大值的城市,但如果 2 个城市在 2 个不同的日期具有相似的最大值,它不会两次返回同一个城市(因为被键:城市减少)。

我希望它足够清楚,任何精度请询问! 非常感谢!

【问题讨论】:

    标签: python apache-spark pyspark rdd


    【解决方案1】:

    要使所有城市的值都等于最大值,您仍然可以使用reduceByKey,但在数组上而不是在值上:

    • 您将行转换为键/值,其中值是元组数组而不是元组
    • 你通过键减少,如果数组包含相同的值,则合并数组,否则保留具有最大值的数组,reduceByKey
    • 您将值数组展平,将键与它们合并,使用flatMap
    • 最后你执行你的排序

    完整的代码如下:

    def merge(array1, array2):
        if array1[0][2] > array2[0][2]:
            return array1
        elif array1[0][2] == array2[0][2]:
            return array1 + array2
        else:
            return array2
    
    
    res = rdd2.map(lambda x: (x[0][0], [(x[0][1], x[0][2], x[1])]))
    rdd3 = res.reduceByKey(lambda x1, x2: merge(x1, x2))
    rdd4 = rdd3.flatMap(lambda x: map(lambda y: (x[0], y[1], y[0], y[2]), x[1]))
    rdd5 = rdd4.sortBy(lambda a: (-a[3], a[2], a[0]))
    

    然后你就可以打印你的RDD了:

    [print(', '.join([row[0], row[1], row[2], str(row[3])])) for row in rdd5.collect()]
    

    您的输入将为您提供以下输出:

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

    【讨论】:

    • 太棒了@Vincent Doba!最后两件事:结果显示为 "City4, 2020-03-27, x4, 5" 而不是 "City4, X4, 2020-03-27, 5"。订单可以通过 reduceByKey。一直在玩 flatMap 顺序(x[0] -> x[1] 等),但结果没有改变,所以我怀疑合并函数是顺序不正确的地方?
    • 另外,输出带有括号(元组),形式为:(City4,X4,2020-03-27,5),如何删除括号?我尝试了并行化但不起作用。
    • @JohnDoe34 我对结果中的行重新排序。你是对的,你必须玩 flatMap 订单。对于元组问题,我需要一些精度:您期望输出什么?连接所有字段的字符串?因为 rdd 中只能有一种类型,字符串、值、元组、对象或数组。
    【解决方案2】:

    您可以使用 Dataframes 工作/输出吗?

    List = [(('City1', '2020-03-27', 'X1'), 44),
            (('City1', '2020-03-28', 'X1'), 44),
            (('City3', '2020-03-28', 'X3'), 15),
            (('City4', '2020-03-27', 'X4'), 5),
            (('City4', '2020-03-26', 'X4'), 4),
            (('City2', '2020-03-26', 'X2'), 14),
            (('City2', '2020-03-25', 'X2'), 4),
            (('City4', '2020-03-25', 'X4'), 1),
            (('City1', '2020-03-29', 'X1'), 1),
            (('City5', '2020-03-25', 'X5'), 15)]
    
    rdd = sc.parallelize(List)
    
    import pyspark.sql.functions as F
    
    df = rdd\
            .toDF()\
            .select('_1.*', F.col('_2').alias('value'))\
            .orderBy(F.desc('value'), F.asc('_2'), F.desc('_1'))
    
    df.show(truncate=False)
    
    +-----+----------+---+-----+
    |_1   |_2        |_3 |value|
    +-----+----------+---+-----+
    |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    |
    |City2|2020-03-25|X2 |4    |
    |City4|2020-03-26|X4 |4    |
    |City4|2020-03-25|X4 |1    |
    |City1|2020-03-29|X1 |1    |
    +-----+----------+---+-----+
    
    

    【讨论】:

    • 非常感谢 Luiz,但我真的需要找到一个不使用数据帧 sql 的解决方案,而是使用带有 RDD 转换和操作的 python。
    【解决方案3】:

    您可以将您的 rdd 转换为 dataframe,然后使用 Spark's window 获取每个城市的最大值,使用此值过滤行,最后根据需要对数据框进行排序:

    from pyspark.sql import functions as F
    from pyspark.sql import Window
    
    window = Window.partitionBy('City').rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
    
    df = rdd.toDF().select(
     F.col('_1._1').alias('city'),
     F.col('_1._2').alias('date'),
     F.col('_1._3').alias('key'),
     F.col('_2').alias('value'),
    ).withColumn('max_value', F.max('value').over(window))\
     .filter(F.col('value') == F.col('max_value'))\
     .drop('max_value')\
     .orderBy(F.desc('value'), F.asc('date'), F.asc('city'))
    

    您会通过输入 rdd 获得以下数据框:

    +-----+----------+---+-----+
    |city |date      |key|value|
    +-----+----------+---+-----+
    |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    |
    +-----+----------+---+-----+
    

    如果您在流程结束时需要一个 RDD,您可以使用.rdd 方法检索它:

    df.rdd
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-08-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多