【问题标题】:pyspark rdd taking the max frequency with the least agepyspark rdd 采用年龄最小的最大频率
【发布时间】:2020-07-06 06:59:17
【问题描述】:

我有一个如下的rdd:

[{'age': 2.18430371791803,
  'code': u'"315.320000"',
  'id': u'"00008RINR"'},
 {'age': 2.80033330216659,
  'code': u'"315.320000"',
  'id': u'"00008RINR"'},
 {'age': 2.8222365762732,
  'code': u'"315.320000"',
  'id': u'"00008RINR"'},
 {...}]

我正在尝试通过使用如下代码获取最高频率代码来将每个 id 减少到仅 1 条记录:

rdd.map(lambda x: (x["id"], [(x["age"], x["code"])]))\
.reduceByKey(lambda x, y: x + y)\
.map(lambda x: [i[1] for i in x[1]])\
.map(lambda x: [max(zip((x.count(i) for i in set(x)), set(x)))])

这个实现有一个问题,它没有考虑年龄,所以如果一个id有多个频率为2的代码,它将采用最后一个代码。

为了说明这个问题,请考虑这个简化的 id:

(u'"000PZ7S2G"',
 [(4.3218651186303, u'"388.400000"'),
  (4.34924421126357, u'"388.400000"'),
  (4.3218651186303, u'"389.900000"'),
  (4.34924421126357, u'"389.900000"'),
  (13.3667102491139, u'"794.310000"'),
  (5.99897016368982, u'"995.300000"'),
  (6.02634923989903, u'"995.300000"'),
  (4.3218651186303, u'"V72.19"'),
  (4.34924421126357, u'"V72.19"'),
  (13.3639723398581, u'"V81.2"'),
  (13.3667102491139, u'"V81.2"')])

我的代码会输出:

[(2, u'"V81.2"')]

当我想让它输出时:

[(2, u'"388.400000"')]

因为虽然这两个代码的频率相同,但代码 388.400000 的年龄较小,出现在最前面。

通过在 .reduceByKey() 之后添加这一行:

.map(lambda x: (x[0], [i for i in x[1] if i[0] == min(x[1])[0]]))

我能够过滤掉那些年龄大于最小年龄的人,但是我只考虑那些年龄最小的人,而不是所有代码来计算它们的频率。我不能在 [max(zip((x.count(i) for i in set(x)), set(x)))] 之后应用相同/相似的逻辑,因为 set(x) 是 x 的集合[1],不考虑年龄。

我应该补充一下,我不想只取频率最高的第一个代码,我想取频率最高、年龄最小的代码,或者首先出现的代码,如果可能的话,仅使用 rdd 操作。

我想要得到的 SQL 中的等效代码类似于:

SELECT code, count(*) as code_frequency
FROM (SELECT id, code, age
FROM (SELECT id, code, MIN(age) AS age, COUNT(*) as cnt,
             ROW_NUMBER() OVER (PARTITION BY id ORDER BY COUNT(*) DESC, MIN(age)) as seqnum
      FROM tbl
      GROUP BY id, code
     ) t
WHERE seqnum = 1) a
GROUP BY code
ORDER by code_frequency DESC
LIMIT 5;

作为 DF(尽管试图避免这种情况):

wc = Window().partitionBy("id", "code").orderBy("age")
wc2 = Window().partitionBy("id")
df = rdd.toDF()
df = df.withColumn("count", F.count("code").over(wc))\
.withColumn("max", F.max("count").over(wc2))\
.filter("count = max")\
.groupBy("id").agg(F.first("age").alias("age"),
                           F.first("code").alias("code"))\
.orderBy("id")\
.groupBy("code")\
.count()\
.orderBy("count", ascending = False)

非常感谢您对此的任何帮助。

【问题讨论】:

  • 这看起来像 How to select the first row of each group? 的副本
  • 这是一个类似的问题,有非常不同的解决方案,没有一个是在 python 中使用 rdd 操作。
  • @mad-a,你的问题有点令人困惑。在等效的 SQL 代码中,计算了 MIN(age),但从未在任何逻辑中使用。最后的聚合听起来合并ids 之间的计算。如果您可以提供具有至少一个不同id 和预期结果的样本数据,将会很有帮助。
  • @jxc 嗨,很抱歉,我忘了在 id 分区中的 order by 语句中的 count(*) desc 之后添加 min(age)。在本质上;如果有 2 个代码具有相同的计数,我希望将 seqnum = 1 分配给最小年龄最小的行。

标签: apache-spark pyspark count rdd reduce


【解决方案1】:

基于您代码的SQL等效,我将逻辑转换为以下rdd1加上一些后处理(从原始RDD开始):

rdd = sc.parallelize([{'age': 4.3218651186303, 'code': '"388.400000"', 'id': '"000PZ7S2G"'},
 {'age': 4.34924421126357, 'code': '"388.400000"', 'id': '"000PZ7S2G"'},
 {'age': 4.3218651186303, 'code': '"389.900000"', 'id': '"000PZ7S2G"'},
 {'age': 4.34924421126357, 'code': '"389.900000"', 'id': '"000PZ7S2G"'},
 {'age': 13.3667102491139, 'code': '"794.310000"', 'id': '"000PZ7S2G"'},
 {'age': 5.99897016368982, 'code': '"995.300000"', 'id': '"000PZ7S2G"'},
 {'age': 6.02634923989903, 'code': '"995.300000"', 'id': '"000PZ7S2G"'},
 {'age': 4.3218651186303, 'code': '"V72.19"', 'id': '"000PZ7S2G"'},
 {'age': 4.34924421126357, 'code': '"V72.19"', 'id': '"000PZ7S2G"'},
 {'age': 13.3639723398581, 'code': '"V81.2"', 'id': '"000PZ7S2G"'},
 {'age': 13.3667102491139, 'code': '"V81.2"', 'id': '"000PZ7S2G"'}])

rdd1 = rdd.map(lambda x: ((x['id'], x['code']),(x['age'], 1))) \
    .reduceByKey(lambda x,y: (min(x[0],y[0]), x[1]+y[1])) \
    .map(lambda x: (x[0][0], (-x[1][1] ,x[1][0], x[0][1]))) \
    .reduceByKey(lambda x,y: x if x < y else y) 
# [('"000PZ7S2G"', (-2, 4.3218651186303, '"388.400000"'))]

地点:

  1. 使用map初始化pair-RDD,key=(x['id'], x['code']), value=(x['age'], 1)
  2. 使用reduceByKey计算min_agecount
  3. 使用map重置pair-RDD,key=id和value=(-count, min_age, code)
  4. 使用reduceByKey 找到相同id 的元组(-count, min_age, code) 的最小值

以上步骤类似:

  • 步骤(1)+(2):groupby('id', 'code').agg(min('age'), count())
  • 步骤(3)+(4):groupby('id').agg(min(struct(negative('count'),'min_age','code')))

然后您可以通过rdd1.map(lambda x: (x[0], x[1][2], x[1][1])) 在您的SQL 中获取派生表a,但这一步不是必需的。 code可以直接从上面的rdd1中通过另外一个map函数+countByKey()方法来统计,然后对结果进行排序:

sorted(rdd1.map(lambda x: (x[1][2],1)).countByKey().items(), key=lambda y: -y[1])
# [('"388.400000"', 1)]

但是,如果您要查找的是所有 ids 的总和(计数),请执行以下操作:

rdd1.map(lambda x: (x[1][2],-x[1][0])).reduceByKey(lambda x,y: x+y).collect()
# [('"388.400000"', 2)]

【讨论】:

  • 非常感谢,你是绝对的英雄!感谢您将其分解并如此简洁地解释它。 :)
【解决方案2】:

如果将 rdd 转换为数据框是一种选择,我认为这种方法可以解决您的问题:

from pyspark.sql.functions import row_number, col
from pyspark.sql import Window
df = rdd.toDF()
w = Window.partitionBy('id').orderBy('age')
df = df.withColumn('row_number', row_number.over(w)).where(col('row_number') == 1).drop('row_number')

【讨论】:

  • 我试图坚持只使用 rdd 动作。但我不确定你的代码是否真的可以工作,我将编辑 op 以包含 df 版本。
猜你喜欢
  • 2018-10-29
  • 1970-01-01
  • 2015-11-02
  • 1970-01-01
  • 1970-01-01
  • 2014-07-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多