为MapReduce输出添加品牌总计数的Reducer实现问题
Hadoop Reducer 修正:为汽车品牌记录添加总出现次数
输入数据格式
chevy, {mod: spark | col: brown} chevy, {mod: equinox | col: red} honda, {mod:civic | col:black} honda, {mod:accord | col:white} honda, {mod:crv | col:pink} honda, {mod:hrv | col:gray} toyota, {mod:corola | col:white}
预期输出
每条记录需追加对应品牌的总出现次数:
chevy, {mod: spark | col: brown | total:2} chevy, {mod: equinox | col: red | total:2} honda, {mod:civic | col:black | total:4} honda, {mod:accord | col:white | total:4} honda, {mod:crv | col:pink | total:4} honda, {mod:hrv | col:gray | total:4} toyota, {mod:corola | col:white | total:1}
现有代码问题
原Reducer代码采用递增计数逻辑,导致每条记录的total值是当前条目在品牌内的顺序号,而非品牌总出现次数。核心问题是没有先统计同一品牌的所有记录数量,再批量输出。
修正后的Reducer代码
import sys def main(): current_brand = None brand_records = [] for line in sys.stdin: line = line.strip() if not line: continue # 拆分品牌和属性部分 brand_part, attr_part = line.split(', ', 1) attr_part = attr_part.strip() if current_brand != brand_part: # 处理上一个品牌的所有记录 if current_brand is not None: total = len(brand_records) for record in brand_records: updated_attr = record.rstrip('}') + f' | total:{total}}}' print(f"{current_brand}, {updated_attr}") # 重置当前品牌和记录列表 current_brand = brand_part brand_records = [attr_part] else: # 同一品牌,添加到记录列表 brand_records.append(attr_part) # 处理最后一个品牌的记录 if current_brand is not None: total = len(brand_records) for record in brand_records: updated_attr = record.rstrip('}') + f' | total:{total}}}' print(f"{current_brand}, {updated_attr}") if __name__ == "__main__": main()
代码说明
- 收集记录:遍历输入时,先将同一品牌的所有属性记录收集到列表中,不立即输出。
- 批量输出:当品牌切换时,计算该品牌的总记录数,然后给每条记录追加
total字段后输出。 - 收尾处理:循环结束后,处理最后一个品牌的剩余记录,避免遗漏。
内容的提问来源于stack exchange,提问作者nkah
相关产品推荐
相关产品推荐

