如何对reduceByKey结果再次执行reduceByKey以分析年度销售长尾效应
解决步骤
1. 先统计各年度每个产品的出现次数
不需要先按年份分组再处理产品ID,直接用(年份, 产品ID)作为复合键,一次reduceByKey就能高效得到结果:
# 生成(年份, 产品ID)作为键,值为1,求和得到各年度各产品的出现次数 product_year_count_rdd = data_rdd.map(lambda x: ((x.Year, x.asin), 1)).reduceByKey(lambda a, b: a + b)
这一步会得到类似((2020, 'asin1'), 5), ((2020, 'asin2'), 3), ((2021, 'asin3'), 10)的RDD,直接拿到每个年度每个产品的计数,比groupByKey更高效(reduceByKey会先在分区内聚合,减少shuffle的数据量)。
2. 按年份分组,判断长尾效应
将同一年份的产品数据聚合后,计算前20%产品的累计占比,判断是否达到65%的阈值:
# 转换键为年份,值为(产品ID, 次数),再按年份分组 year_products_rdd = product_year_count_rdd.map(lambda x: (x[0][0], (x[0][1], x[1]))).groupByKey() # 定义函数判断单年度是否存在长尾效应 def check_long_tail(products): # 按产品出现次数降序排序 sorted_products = sorted(products, key=lambda x: x[1], reverse=True) total_count = sum(p[1] for p in sorted_products) if total_count == 0: return False # 无数据的年份直接返回False # 计算20%产品数量的阈值,避免产品数量过少时出现0的情况 product_num = len(sorted_products) top_20_percent_num = max(int(product_num * 0.2), 1) # 计算前20%产品的累计次数及占比 top_total = sum(p[1] for p in sorted_products[:top_20_percent_num]) return (top_total / total_count) >= 0.65 # 应用函数得到各年度的长尾效应判断结果 year_long_tail_rdd = year_products_rdd.map(lambda x: (x[0], check_long_tail(x[1]))) # 查看最终结果 year_long_tail_rdd.collect()
若坚持初始思路的修改方案
如果一定要先按年份分组再处理产品ID,可以这样调整(性能略低于复合键方式):
# 先按年份分组,得到每个年份对应的产品ID列表 year_asin_list_rdd = data_rdd.map(lambda x: (x.Year, x.asin)).groupByKey() # 在每个年份组内统计产品ID的出现次数 year_product_count_rdd = year_asin_list_rdd.mapValues(lambda asins: {asin: asins.count(asin) for asin in set(asins)} )
后续同样可以基于year_product_count_rdd调用check_long_tail函数完成判断。
内容的提问来源于stack exchange,提问作者user20831785
相关产品推荐
相关产品推荐

