Spark中值为嵌套字典的(Key,Value)元组Reduce聚合实现方法
Spark RDD 嵌套字典值的聚合实现方案
你当前的实现问题出在reduceByKey传入的合并逻辑仅做了字典覆盖,未对同键的数值做累加:{**a,**b}是Python的字典解符合并操作,若两个字典存在相同键,后一个字典的键值会直接覆盖前一个的,因此最终仅保留了最后处理到的分区对应值,没有累加效果。
实现方案
首先自定义字典合并函数,实现相同键数值累加、不同键保留的逻辑,适配你使用的Python 3.7版本:
def merge_dicts(dict_a, dict_b): # 取两个字典所有键的并集,兼容Python3.7写法 all_keys = set(dict_a.keys()).union(dict_b.keys()) return { key: dict_a.get(key, 0) + dict_b.get(key, 0) for key in all_keys }
替换原有reduceByKey的匿名函数即可,修改后的完整代码如下:
errors_distribution = dataFrame.rdd \ .mapPartitions(error_schema_wrapper(validation_schema)) \ .reduceByKey(merge_dicts) \ .collectAsMap()
该逻辑可完全匹配你的需求:针对相同Key对应的字典,会把字典内同字段的数值做累加,字典内独有的字段则直接保留原值。测试你的示例输入后,可准确得到abc对应min:6、def对应min:8、geh对应nullable:1的预期输出。
内容的提问来源于stack exchange,提问作者Jack Miller
相关产品推荐
相关产品推荐

