You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 15:14:51