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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 20:15:00