Python中基于MapReduce的词频频次统计:单作业能否实现?
关于MapReduce实现词频频次统计的作业次数问题
先把需求拆解清楚:咱们要做的不是普通的词频统计,而是在词频统计的基础上,再统计「不同词频值出现的次数」——相当于要完成两级聚合计算。先看你给的示例:
输入单词集合:
(be,to,the,the,now,now,now,see,see,see)
第一步词频统计结果:be:1、to:1、the:2、now:3、see:3
最终目标输出:不同词频的出现次数,即1:2、2:1、3:2
核心结论:必须分两个MapReduce作业完成
MapReduce的单轮作业只能完成一级聚合,因为它的Map->Shuffle->Reduce流程里,Shuffle阶段只会按照一组Key做数据分组分发,没办法同时完成两级分组聚合。咱们的需求刚好是两层逻辑,所以必须拆成两个独立作业:
第一个作业:完成基础词频统计
- Map阶段:读取输入的单词数据,输出键值对
(单词, 1),比如读到the就输出("the", 1) - Shuffle阶段:框架自动把相同单词的键值对分配到同一个Reduce任务
- Reduce阶段:对同一个单词的所有
1求和,输出键值对(单词, 词频),比如("the", 2)
第二个作业:统计词频的出现次数
- Map阶段:读取第一个作业的输出,把键值对反转,输出
(词频, 1),比如读到("the", 2)就输出(2, 1) - Shuffle阶段:框架把相同词频值的键值对分配到同一个Reduce任务
- Reduce阶段:对同一个词频的所有
1求和,输出键值对(词频, 出现次数),比如(1, 2),这就是最终需要的结果
Python实现小提示
如果用Hadoop Streaming这类工具在Python里实现,每个作业可以写两个独立脚本(mapper.py和reducer.py),先提交第一个作业生成词频结果,再把这个结果目录作为第二个作业的输入提交即可。
内容的提问来源于stack exchange,提问作者Kostas Papastamos
相关产品推荐
相关产品推荐

