如何从Spark RDD的map、reduceByKey操作结果中移除None对应的行
操作方案
你可以通过添加filter算子实现过滤,优先选择在聚合前过滤无效数据,减少shuffle阶段的数据传输量,计算效率更高:
方案1:聚合前过滤(推荐)
直接在positive函数映射后就过滤掉无效值,300多万条None对应的记录不需要参与后续聚合,性能收益明显:
fileRDD.map(positive)\ # 仅保留0-6的有效值,自动排除None .filter(lambda x: x in {0,1,2,3,4,5,6})\ .map(lambda x: [x,1])\ .reduceByKey(lambda x,y: x+y)\ .take(10)
如果你只需要排除None,也可以把过滤条件简化为.filter(lambda x: x is not None)。
方案2:聚合后过滤
如果已经完成聚合,只想对输出结果做过滤,可以在reduceByKey之后添加过滤逻辑:
fileRDD.map(positive)\ .map(lambda x: [x,1])\ .reduceByKey(lambda x,y: x+y)\ # 过滤掉键为None的结果 .filter(lambda item: item[0] is not None)\ .take(10)
内容的提问来源于stack exchange,提问作者merchmallow
相关产品推荐
相关产品推荐

