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

如何通过PySpark CountVectorizer提取每行TopN关键词

按行提取TopN高频关键词实现方案

核心实现思路

  • CountVectorizer训练完成后,模型的vocabulary属性存储了词表列表,列表索引和输出稀疏向量的位置索引完全对应,比如vocabulary[0]就是稀疏向量第0位代表的词汇。
  • 稀疏向量结构为(总词数, [出现过的词索引], [对应词的词频]),同一位置的索引和词频一一对应。
  • 自定义函数逐行处理稀疏向量:将索引和对应词频配对后按词频降序排序,取前N个索引映射为原始词汇,拼接为目标字符串即可。
  • 注意原有拆分逻辑的问题:直接用逗号拆分文本会导致非首词带前导空格(比如拆分后得到 mars而非mars),需要加trim操作清理空格,避免同一个词被识别为两个不同词。

完整可运行代码示例

from pyspark.ml.feature import CountVectorizer
from pyspark.sql.functions import split, col, udf, expr
from pyspark.sql.types import StringType

# 1. 文本拆分(新增trim清理每个词的前后空格)
input_df = input_df.withColumn("text_array", expr("transform(split(text, ','), x -> trim(x))"))

# 2. 原有CountVectorizer训练逻辑
cv_text = CountVectorizer() \
    .setInputCol("text_array") \
    .setOutputCol("cv_text")
cv_model = cv_text.fit(input_df)
cv_result = cv_model.transform(input_df)

# 3. 配置TopN参数,获取词表
TOP_N = 2
vocab_list = cv_model.vocabulary

# 4. 定义提取TopN关键词的UDF
def extract_topn(vector):
    # 配对(词频, 词索引),按词频降序排序
    freq_idx_pairs = sorted(zip(vector.values, vector.indices), reverse=True)
    # 取前N个索引映射为词汇
    topn_words = [vocab_list[idx] for _, idx in freq_idx_pairs[:TOP_N]]
    # 拼接为逗号分隔字符串
    return ", ".join(topn_words)

# 注册UDF并应用到数据集
extract_udf = udf(extract_topn, StringType())
final_df = cv_result.withColumn("keywords", extract_udf(col("cv_text")))

# 查看结果
final_df.show(truncate=False)

补充说明

  • 如果需要自定义词频相同时的排序规则,可以修改sorted的key参数,比如加入词在词表中的顺序作为次排序维度,保证结果稳定。
  • 如果需要输出关键词数组而非拼接字符串,只要把UDF的返回类型改成ArrayType(StringType()),去掉最后一步join操作即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 02:18:34