【发布时间】: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)
非常感谢您对此的任何帮助。
【问题讨论】:
-
这是一个类似的问题,有非常不同的解决方案,没有一个是在 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