Pyspark如何修改代码实现列表数字出现频率统计
PySpark频次统计修改方案
你可以直接使用PySpark RDD内置的countByValue()方法实现频次统计,完整修改后代码如下:
# 第一步:统计每个值的出现频次 freq_dict = ARDD.map(function_B) \ .filter(lambda x: x is not None) \ .map(lambda x: int(x)) # 把字符串类型的数字转成整数,不需要可以删掉这行 .countByValue() # 第二步:按频次降序、数值升序排序,拼接成要求的格式 sorted_items = sorted(freq_dict.items(), key=lambda item: (-item[1], item[0])) result = [f"{k}:{v}" for k, v in sorted_items] # 如果需要直接输出你示例的无引号格式,执行下面的打印语句 print("[" + ", ".join(result) + "]")
大数量场景可选优化方案
如果你的数据量极大,不想把全量统计结果拉取到Driver端,可以用reduceByKey算子替换,代码如下:
freq_rdd = ARDD.map(function_B) \ .filter(lambda x: x is not None) \ .map(lambda x: (int(x), 1)) \ .reduceByKey(lambda a, b: a + b) \ .sortBy(lambda x: (-x[1], x[0])) # 最终输出 result = [f"{k}:{v}" for k, v in freq_rdd.collect()] print("[" + ", ".join(result) + "]")
以上两种方案执行后输出结果均为:
[2:2, 3:2, 10:1, 12:1]
内容的提问来源于stack exchange,提问作者pouchewar
相关产品推荐
相关产品推荐

