Scala中统计Mongo查询结果文档指定字段的单词出现次数
错误原因
你遇到的报错有两个核心问题:
reduceByKey是Spark RDD的专属算子,collection.find()返回的是Mongo驱动的FindIterable类型,本身没有这个方法,直接调用必然报符号未解析错误。- 你之前的
map逻辑错误,直接将整个文档对象当做key,没有提取文本字段分词,也没有对应你要统计的目标单词,自然会触发+符号的解析错误。
实现方案
分两种场景给出实现代码:
场景1:小数据量,纯Scala本地处理(不用Spark)
如果你处理的数据量不大,不需要分布式计算,直接遍历查询结果统计即可,这里以Mongo同步Scala驱动为例:
import org.mongodb.scala.model.Filters._ import scala.collection.mutable.Map // 1. 定义你要统计的目标单词集合,不需要区分大小写就统一转小写 val targetWords = Set("boulder", "denver", "cowx") // 初始化计数容器 val wordCount = Map[String, Int]().withDefaultValue(0) // 2. 执行你的原有查询 val f = collection.find( and( geoWithinCenter("geo.coordinates", lon, lat, radius), gt("EpochTime", start), lt("EpochTime", end) ) ) // 3. 遍历结果统计 val docList = f.toList() docList.foreach { doc => // 注意:你给出的样例中实际文本字段为TweetText,如果你实际字段名是text请修改此处 val text = doc.getString("TweetText").toLowerCase // 按非字母数字符号分词,可根据需求调整正则规则 val words = text.split("\\W+") words.foreach { word => if (targetWords.contains(word)) { wordCount(word) += 1 } } } // 输出统计结果 wordCount.foreach { case (word, count) => println(s"单词$word 出现次数:$count") }
场景2:大数据量,用Spark DataFrame实现
如果你要处理的数据量较大,用Mongo Spark Connector直接读成DataFrame处理即可:
import org.apache.spark.sql.SparkSession import com.mongodb.spark._ import org.apache.spark.sql.functions._ // 初始化SparkSession,替换为你自己的Mongo连接地址 val spark = SparkSession.builder() .appName("MongoWordCount") .config("spark.mongodb.input.uri", "mongodb://<你的Mongo地址>/<库名>.<集合名>") .getOrCreate() // 1. 读取数据并按你的条件过滤 val df = MongoSpark.load(spark) .filter( expr(s"geoWithinCenter(geo.coordinates, $lon, $lat, $radius)") && col("EpochTime") > start && col("EpochTime") < end ) // 2. 定义要统计的目标单词 val targetWords = Seq("boulder", "denver", "cowx") // 3. 分词、过滤、统计 val wordCountDf = df .select(lower(col("TweetText")).alias("text")) .select(split(col("text"), "\\W+").alias("words")) .select(explode(col("words")).alias("word")) .filter(col("word").isin(targetWords: _*)) .groupBy("word") .count() // 输出结果 wordCountDf.show()
注意事项
- 请根据你实际的文本存储字段名调整代码中的字段,你问题中描述的是
text字段,但给出的样例JSON中文本存储在TweetText字段。 - 如果你需要统计所有单词的出现次数,直接去掉代码中过滤目标单词的步骤即可。
- 分词规则可根据你的需求调整,比如要区分大小写就去掉
toLowerCase的调用,要保留特殊符号就修改split的正则表达式。
内容的提问来源于stack exchange,提问作者Al-Meqdad Jabi
相关产品推荐
相关产品推荐

