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

Spark实现按累计占比阈值输出排序的高频词统计结果

Solution for Filtering Top Words by Cumulative Frequency Percentage

Got it, let's walk through how to finish your Spark code to get top words until their cumulative count hits your target percentage. You've already laid the groundwork with word counting—now we just need to add the sorting and cumulative filtering logic.

Step 1: Full RDD-based Implementation

First, let's extend your existing code with straightforward RDD operations. Note that we'll pull the sorted results to the driver (which is safe for top words since there won't be millions of them):

// 将文件拆分为单词
val lines = sc.textFile("somefile.txt")
val words = lines.flatMap(line => line.split(" "))
words.persist()
val wordCount = words.count()
val wordCounts = words.map(word => (word, 1)).reduceByKey{case (x, y) => x + y}

// 1. Sort word counts in descending order of frequency
val sortedWordCounts = wordCounts.sortBy(-_._2)

// 2. Set your target percentage (e.g., 80% of total words)
val targetPercentage = 80.0

// 3. Collect sorted results to driver and filter by cumulative percentage
val sortedWordList = sortedWordCounts.collect()
var cumulativeCount = 0L
val filteredTopWords = sortedWordList.takeWhile { case (word, count) =>
  cumulativeCount += count
  (cumulativeCount.toDouble / wordCount) * 100 <= targetPercentage
}

// 4. Print the final results
filteredTopWords.foreach { case (word, count) =>
  val currentCumulativePct = (cumulativeCount.toDouble / wordCount) * 100
  println(s"Word: *$word*, Count: $count, Cumulative Percentage: ${currentCumulativePct.round}%")
}

Step 2: Alternative with Spark SQL (Better for Large Datasets)

If you're working with extremely large datasets where even top words might be too many to collect to the driver, using Spark SQL window functions is a better distributed approach:

import spark.implicits._

// Reuse your existing wordCounts RDD
val wordCountDF = wordCounts.toDF("word", "count")

// Define a window to sort by frequency and calculate cumulative sums
val windowSpec = Window.orderBy($"count".desc)
val cumulativeDF = wordCountDF
  .withColumn("cumulative_count", sum($"count").over(windowSpec))
  .withColumn("cumulative_percentage", ($"cumulative_count" / wordCount) * 100)
  .filter($"cumulative_percentage" <= targetPercentage)

// Show the results (or write to output)
cumulativeDF.select("word", "count", "cumulative_percentage").show()

Key Notes

  • The RDD approach is simple and fast for most cases—since we're only dealing with top words, collecting them to the driver won't cause memory issues.
  • The Spark SQL method handles distributed cumulative calculations, so it's safer for massive datasets where you can't pull all top words to the driver.
  • Adjust targetPercentage to your desired threshold (e.g., 50 for 50%, 90 for 90%).

内容的提问来源于stack exchange,提问作者Bassinator

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:41:18