如何通过PySpark HashingTF统计包含指定词条的文档数量
问题解决方案
为什么HashingTF不适合当前需求
HashingTF使用哈希Trick完成词到向量索引的映射,映射过程不可逆,无法从索引反向还原对应的原始词,同时存在哈希碰撞风险(多个词映射到同一个索引),会导致统计结果失真,因此不推荐用HashingTF实现该需求。
最优实现方案(原生Spark DataFrame API)
该方案直接基于分词后的数组做计算,逻辑直观、效率更高,无结果失真问题。
前置预处理
首先统一文本大小写避免统计误差,同时给每行文档添加唯一ID方便后续统计:
from pyspark.sql.functions import split, col, lower, monotonically_increasing_id # 统一转小写、添加唯一文档ID、分词(用\\s+匹配任意空白字符更鲁棒) spark_df = spark_df.withColumn("doc_id", monotonically_increasing_id()) \ .withColumn("words", split(lower(col("text")), "\\s+"))
1. 统计包含单词lazy的文档总行数
直接用array_contains判断分词数组中是否存在目标词后计数即可:
from pyspark.sql.functions import array_contains cnt_lazy = spark_df.filter(array_contains(col("words"), "lazy")).count()
2. 统计同时包含dog和bird的文档行数
组合多个array_contains条件过滤:
cnt_dog_bird = spark_df.filter( array_contains(col("words"), "dog") & array_contains(col("words"), "bird") ).count()
3. 全量词文档频率统计、最高频词及对应文档获取
如果需要统计所有词的出现文档数、找最高频词等扩展需求,可以通过拆分词后分组统计实现:
from pyspark.sql.functions import explode, countDistinct # 计算每个词的文档出现次数 doc_freq_df = spark_df.select("doc_id", explode(col("words")).alias("term")) \ .dropDuplicates(["doc_id", "term"]) # 同一个文档同一个词只算一次 .groupBy("term") \ .agg(countDistinct("doc_id").alias("doc_count")) \ .orderBy(col("doc_count").desc()) # 获取最高频词 top_term = doc_freq_df.first() print(f"最高频词:{top_term['term']},出现文档数:{top_term['doc_count']}") # 获取包含最高频词的所有文档 top_term_docs = spark_df.filter(array_contains(col("words"), top_term["term"]))
(不推荐)基于HashingTF的实现方案
如果一定要用已生成的HashingTF结果做统计,需要先计算目标词对应的哈希索引,再判断向量对应位置的值是否大于0,该方案存在哈希碰撞误差风险:
def get_term_index(term, num_features=262144): # 和HashingTF默认的哈希逻辑对齐 return hash(term) % num_features # 统计包含lazy的文档数 lazy_index = get_term_index("lazy") cnt_lazy_hashing = df_term_frequency.filter(col("features")[lazy_index] > 0).count()
内容的提问来源于stack exchange,提问作者sbs0202
相关产品推荐
相关产品推荐

