PySpark如何按语言分组统计词频并生成词频映射字典
实现方案
你需要先统计单词汇的出现次数,再聚合为映射结构,完整可运行代码如下:
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 原有构造数据的逻辑 spark = SparkSession.builder.appName("word_count_map").getOrCreate() sdf = spark.createDataFrame([ ('eng', "cat"), ('eng', 'cat'), ('eng','dog'), ('eng','cat') ], ["lang", "text"]) # 核心实现逻辑 # 1. 按语言+词汇双字段分组,统计每个词汇的单独出现次数 word_count_per_lang = sdf.groupBy("lang", "text").agg(F.count("*").alias("cnt")) # 2. 仅按语言分组,拼接词汇和计数为Map结构 result = word_count_per_lang.groupBy("lang").agg( F.map_from_entries( F.collect_list(F.create_map(F.col("text"), F.col("cnt"))) ).alias("count_words") ) # 输出结果 result.show(truncate=False)
运行后输出效果和需求完全匹配:
+----+------------------+ |lang|count_words | +----+------------------+ |eng |{cat -> 3, dog -> 1}| +----+------------------+
逻辑说明
groupBy("lang", "text")双字段分组,确保统计的是每个语言下每个词汇的单独计数create_map(text, cnt)将每行的词汇和对应计数生成为单键值对的Map结构collect_list()将同语言下的所有单键值对Map收集为数组map_from_entries()将键值对数组合并为完整的Map结构,即你需要的字典格式
如果需要将输出的Spark Map格式转为Python原生字典,可直接对查询结果的行数据调用asDict()方法处理。
内容的提问来源于stack exchange,提问作者Rory
相关产品推荐
相关产品推荐

