【问题标题】:How to execute this sorting process in pyspark?如何在 pyspark 中执行这个排序过程?
【发布时间】:2017-04-08 21:19:16
【问题描述】:

我尝试了 map、mapValues 和 sort,但没有任何效果。 问题描述如下: “根据相似度(值中的第二个),如果相同,则选择ID最小的用户(值中的第一个)。” 键值对的列表是:

[
    (18, [(2, 0.5)]),
    (30, [(19, 0.5), (6, 0.25)]),
    (6, [(30, 0.25), (20, 0.2), (19, 0.2)]),
    (19, [(30, 0.5), (8, 0.2), (6, 0.2)]),
    (2, [(18, 0.5)]),
    (26, [(9, 0.2)]),
    (9, [(26, 0.2)])
]

我想得到:

[
    (18, [(2, 0.5)]),
    (30, [(19, 0.5), (6, 0.25)]),
    (6, [(30, 0.25), (19, 0.2)]),
    (19, [(30, 0.5), (6, 0.2)]),
    (2, [(18, 0.5)]),
    (26, [(9, 0.2)]),
    (9, [(26, 0.2)])
]

非常感谢!

【问题讨论】:

    标签: sorting pyspark


    【解决方案1】:

    非常简单。只需要弄清楚必要的转换:

    data = [(18, [(2, 0.5)]),
     (30, [(19, 0.5), (6, 0.25)]),
     (6, [(30, 0.25), (20, 0.2), (19, 0.2)]),
     (19, [(30, 0.5), (8, 0.2), (6, 0.2)]),
     (2, [(18, 0.5)]),
     (26, [(9, 0.2)]),
     (9, [(26, 0.2)])]
    
    rdd1 = sc.parallelize(data)
    
    rdd2 = rdd1.flatMapValues(lambda x:x)
    
    rdd3 = rdd2.map(lambda x: ((x[0], x[1][1]),x[1][0]))
    
    rdd4 = rdd3.reduceByKey(min)
    
    rdd5 = rdd4.map(lambda x: (x[0][0], (x[1], x[0][1])))
    
    rdd6 = rdd5.reduceByKey(lambda x,y: [x,y])
    rdd6.collect()
    [(9, (26, 0.2)),
     (26, (9, 0.2)),
     (18, (2, 0.5)),
     (30, [(6, 0.25), (19, 0.5)]),
     (2, (18, 0.5)),
     (6, [(30, 0.25), (19, 0.2)]),
     (19, [(30, 0.5), (6, 0.2)])]
    

    【讨论】:

      猜你喜欢
      • 2020-06-06
      • 1970-01-01
      • 2011-09-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-03-21
      • 2014-07-11
      • 2015-06-26
      相关资源
      最近更新 更多