Scala中DStream的RDD重复词共现频率统计问题
单词共现频率统计问题(Scala + DStream)
问题场景
我是Scala新手,在统计DStream中每个RDD的单词共现频率时遇到以下问题:
- 输入无重复单词时,代码运行正常
- 输入存在重复单词时,重复词的共现会被跳过
- 移除过滤逻辑后,会错误统计单词自身的共现情况,不符合需求
输入示例
输入内容:like pig like hive
原有代码及问题
带过滤的代码
val coOccurrenceCountsB = rdd.flatMap { line => val words = line.split("\\s+").filter(word => word.matches("[a-zA-Z]+") && word.length >= 3) words.flatMap { word1 => words.filter(word2 => word2 != word1).map(word2 => ((word1, word2), 1)) } }.reduceByKey(_ + _)
带过滤的输出
((pig,hive),1) ((hive,like),2) ((like,hive),2) ((pig,like),2) ((hive,pig),1) ((like,pig),2)
问题:重复词like的共现情况完全被过滤掉,未出现在结果中。
无过滤的输出(移除filter(word2 => word2 != word1)后)
((pig,hive),1) ((pig,pig),1) ((hive,like),2) ((like,hive),2) ((pig,like),2) ((hive,pig),1) ((like,like),4) ((like,pig),2) ((hive,hive),1)
问题:出现(pig,pig)、(like,like)这类单词自身的共现统计,且(like,like)计数错误(实际应为2,此处算成4)。
期望输出
((pig,hive),1) ((hive,like),2) ((like,hive),2) ((pig,like),2) ((hive,pig),1) ((like,pig),2) ((like, like), 2)
修正方案及代码
核心思路
原代码的问题在于通过单词内容判断是否过滤:要么过滤掉所有内容相同的词(导致重复词共现丢失),要么允许自身配对(导致错误统计)。
正确做法是通过单词在原行的索引过滤:给每个单词绑定它在原行的位置索引,只配对索引不同的单词——既保留内容相同但位置不同的单词共现,又避免同一位置的单词和自身配对。
修正后的代码
val coOccurrenceCountsB = rdd.flatMap { line => val words = line.split("\\s+").filter(word => word.matches("[a-zA-Z]+") && word.length >= 3) val indexed = words.zipWithIndex // 给每个单词绑定索引 indexed.flatMap { case (word1, index1) => indexed .filter { case (_, index2) => index1 != index2 } // 过滤索引相同的自身配对 .map { case (word2, _) => ((word1, word2), 1) } } }.reduceByKey(_ + _)
修正后效果
运行该代码可得到符合预期的输出,正确统计重复词的共现频率,同时避免单词自身的无效共现。
内容的提问来源于stack exchange,提问作者Fay
相关产品推荐
相关产品推荐

