PySpark计算TF-IDF


tf-idf是一种用于信息检索与文本挖掘的常用加权技术。tf-idf是一种统计方法,用以评估一字词对于一个文档集或一个语料库中的其中一份文档的重要程度。字词的重要性随着它在文件中出现的次数成正比增加,但同时会随着它在语料库中出现的频率成反比下降。

1. TF

在一份给定的文档d dd里,词频(term frequency,tf)指的是某一个给定的词语在该文档中出现的频率。这个数字是对词数(term count)的归一化,以防止它偏向长的文档。(同一个词语在长文档里可能会比短文档有更高的词数,而不管该词语重要与否。)对于在某一特定文档里的词语w i 来说,它的重要性可表示为:

2. IDF

逆向文档频率(inverse document frequency,idf)是一个词语普遍重要性的度量。某一特定词语的idf,可以由总文档数目除以包含该词语之文档的数目,再将得到的商取以10为底的对数得到:

如果一个词越常见,那么分母就越大,逆文档频率(IDF)就越小越接近0。分母之所以要加1,是为了避免分母为0(即所有文档都不包含该词),log表示对得到的值取对数。

3. TF-IDF

某一特定文档内的高词语频率,以及该词语在整个文档集合中的低文档频率,可以产生出高权重的tf-idf。因此,tf-idf倾向于过滤掉常见的词语,保留重要的词语。
T F ? I D F = T F ? I D F 
原文链接:https://blog.csdn.net/mingo220/article/details/106341172

4. PySpark计算TF-IDF

from pyspark.sql import SparkSession

from offline import SparkSessionBase

class TestSparkCountVectorizer(SparkSessionBase):
    def __init__(self):
        self.spark = self._create_spark_session()

    def _create_DataFrame(self):
        df = self.spark.createDataFrame([(1, 'T really liked this movie'),
                      (2, 'I would recommend this movie to my friends'),
                      (3, 'movie was alright but acting was horrible'),
                      (4, 'I am never watching that movie ever again i liked it')],
                     ['user_id', 'review'])
        return df


def TokenizerTest():
    """
    标记化
    :return:
    """
    from pyspark.ml.feature import Tokenizer

    oa = TestSparkCountVectorizer()
    df = oa._create_DataFrame()
    df.show()

    tokenization = Tokenizer(inputCol='review',outputCol='tokens')
    tokenization_df = tokenization.transform(df)
    tokenization_df.show()
    """
            +-------+--------------------+--------------------+
        |user_id|              review|              tokens|
        +-------+--------------------+--------------------+
        |      1|T really liked th...|[t, really, liked...|
        |      2|I would recommend...|[i, would, recomm...|
        |      3|movie was alright...|[movie, was, alri...|
        |      4|I am never watchi...|[i, am, never, wa...|
        +-------+--------------------+--------------------+
    """
    return tokenization_df

def TestStopWordsRemover(tokenization_df):
    """
    停用词移除
    :param df:
    :return:
    """

    from pyspark.ml.feature import StopWordsRemover

    stopword_removal = StopWordsRemover(inputCol='tokens', outputCol='refind_tokens')
    refind_df = stopword_removal.transform(tokenization_df)
    refind_df.select(['user_id', 'tokens', 'refind_tokens']).show(4, False)
    """
            +-------+----------------------------------------------------------------+-------------------------------------+
        |user_id|tokens                                                          |refind_tokens                        |
        +-------+----------------------------------------------------------------+-------------------------------------+
        |1      |[t, really, liked, this, movie]                                 |[really, liked, movie]               |
        |2      |[i, would, recommend, this, movie, to, my, friends]             |[recommend, movie, friends]          |
        |3      |[movie, was, alright, but, acting, was, horrible]               |[movie, alright, acting, horrible]   |
        |4      |[i, am, never, watching, that, movie, ever, again, i, liked, it]|[never, watching, movie, ever, liked]|
        +-------+----------------------------------------------------------------+-------------------------------------+
    """
    return refind_df
def TestCountVectorizer(refind_df):
    """
    ####词袋   数值形式表示文本数据
    ###计数向量器  会统计特定文档中单词出现的次数
    :return:
    """
    from pyspark.ml.feature  import CountVectorizer

    count_vec = CountVectorizer(inputCol='refind_tokens', outputCol='features')
    # 训练词频统计模型
    cv_model = count_vec.fit(refind_df)
    # cv_model.write().overwrite().save("hdfs://hadoop:9000/headlines/models/TestCV.model")

    # 将词频统计模型转换为 DataFrame
    cv_df = cv_model.transform(refind_df)
    cv_df.select(['user_id', 'refind_tokens', 'features']).show(4, False)
    """
            +-------+-------------------------------------+--------------------------------------+
        |user_id|refind_tokens                        |features                              |
        +-------+-------------------------------------+--------------------------------------+
        |1      |[really, liked, movie]               |(11,[0,1,4],[1.0,1.0,1.0])            |
        |2      |[recommend, movie, friends]          |(11,[0,3,8],[1.0,1.0,1.0])            |
        |3      |[movie, alright, acting, horrible]   |(11,[0,2,5,10],[1.0,1.0,1.0,1.0])     |
        |4      |[never, watching, movie, ever, liked]|(11,[0,1,6,7,9],[1.0,1.0,1.0,1.0,1.0])|
        +-------+-------------------------------------+--------------------------------------+
    其中11是单词的个数,可以看到movie出现的频率最高,所以排在最前面。[0,1,7]分别表示单词位于第0个、第1个和第7个位置,[1.0,1.0,1.0]表示单词在本文档中出现的次数
    """
    return cv_df

def TestCountVectorizerModel(refind_df):
    """
    训练idf模型
    :return:
    """
    # 1: 获取词频统计模型
    from pyspark.ml.feature import CountVectorizerModel
    cv_model = CountVectorizerModel.load("hdfs://hadoop:9000/headlines/models/TestCV.model")
    # 2: 获取词频结果
    cv_result = cv_model.transform(refind_df)
    print("-------------------------------词频结果-------------------------")
    cv_result.show()
    """
            -------------------------------词频结果-------------------------
        +-------+--------------------+--------------------+--------------------+--------------------+
        |user_id|              review|              tokens|       refind_tokens|            features|
        +-------+--------------------+--------------------+--------------------+--------------------+
        |      1|T really liked th...|[t, really, liked...|[really, liked, m...|(11,[0,1,4],[1.0,...|
        |      2|I would recommend...|[i, would, recomm...|[recommend, movie...|(11,[0,3,8],[1.0,...|
        |      3|movie was alright...|[movie, was, alri...|[movie, alright, ...|(11,[0,2,5,10],[1...|
        |      4|I am never watchi...|[i, am, never, wa...|[never, watching,...|(11,[0,1,6,7,9],[...|
        +-------+--------------------+--------------------+--------------------+--------------------+

    """
    # 3: 训练IDF模型
    from pyspark.ml.feature import IDF
    idf = IDF(inputCol='features',outputCol='idfFeatures')
    idfModel = idf.fit(cv_result)
    idfModel.write().overwrite().save("hdfs://hadoop:9000/headlines/models/TestIDF.model")
    print(idfModel.idf.toArray()[:20])
    """
        [0.         0.51082562 0.91629073 0.91629073 0.91629073 0.91629073
     0.91629073 0.91629073 0.91629073 0.91629073 0.91629073]
    """


if __name__ == '__main__':
    tokenization_df = TokenizerTest()
    refind_df = TestStopWordsRemover(tokenization_df)
    cv_df = TestCountVectorizer(refind_df)
    TestCountVectorizerModel(refind_df)