spark之CountVectorizer


CountVectorizer会统计特定文档中单词出现的次数,并且会根据单词的频率进行排序,频率高的排在前面,当频率相同时,则它的位置个人感觉是随机的。因为太过例子跑出来,每一次都不相同。

1: 初始化spark环境

# 初始化spark信息
import os
import sys
BASE_DIR= os.path.dirname(os.path.dirname("/bigdata/projects/toutiao_projects/reco_sys/offline/full_cal"))
sys.path.insert(0,os.path.join(BASE_DIR))
PYSPARK_PYTHON = "/usr/local/python3/bin/python3"
# 当存在多个版本时,不指定很可能会导致出错
os.environ["PYSPARK_PYTHON"] = PYSPARK_PYTHON
os.environ["PYSPARK_DRIVER_PYTHON"] = PYSPARK_PYTHON

2:初始化spark配置

"""
离线相关的计算
spark 初始化相关配置
"""
from  pyspark.sql import SparkSession
from pyspark import SparkConf
from pyspark.sql import HiveContext
import os

class SparkSessionBase(object):
    SPARK_APP_NAME = None
    SPARK_URL = 'local'
    SPARK_EXECUTOR_MEMORY = '2g'
    SPARK_EXECUTOR_CORE = 2
    SPARK_EXECUTOR_INSTANCES = 2
    ENABLE_HIVE_SUPPORT = False
    HIVE_URL = "hdfs://hadoop:9000//user/hive/warehouse"
    HIVE_METASTORE = "thrift://hadoop:9083"


    def _create_spark_session(self):
        conf = SparkConf()  # 创建spark config 对象
        config = (
            ("spark.app.name",self.SPARK_APP_NAME), # 设置启动spark 的app名字
            ("spark.executor.memory",self.SPARK_EXECUTOR_MEMORY), # 设置app启动时占用的内存用量
            ('spark.executor.cores',self.SPARK_EXECUTOR_CORE), # 设置spark executor使用的CPU核心数,默认是1核心
            ("spark.master", self.SPARK_URL),  # spark master的地址
            ("spark.executor.instances", self.SPARK_EXECUTOR_INSTANCES),
            ("spark.sql.warehouse.dir",self.HIVE_URL),
            ("hive.metastore.uris",self.HIVE_METASTORE),

        )
        conf.setAll(config)

        # 利用config对象,创建spark session
        if self.ENABLE_HIVE_SUPPORT:
            # return "aaa"
            return SparkSession.builder.config(conf=conf).enableHiveSupport().getOrCreate()
        else:
            # return "bbb"
            return SparkSession.builder.config(conf=conf).getOrCreate()

3: 初始化变量

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

4: 利用pyspark.ml.feature.Tokenizer进行分词

"""
标记化  分词
: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..

 5:利用 pyspark.ml.feature.StopWordsRemover  移除停用次

"""
停用词移除
: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]|
+-------+----------------------------------------------------------------+-------------------------------------+

6:利用pyspark.ml.feature.CountVectorizer 统计词频,训练词频模型,保存词频模型

"""
####词袋   数值形式表示文本数据
###计数向量器  会统计特定文档中单词出现的次数
: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])|
+-------+-------------------------------------+--------------------------------------+

7:利用 pyspark.ml.feature.CountVectorizerModel加载词频模型,获取词频结果,利用pyspark.ml.feature.IDF 训练IDF模型

"""
训练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")
idf_result = idfModel.transform(cv_result)
idf_result.show()
print(cv_model.vocabulary)
print(idfModel.idf.toArray()[:20])
+-------+--------------------+--------------------+--------------------+--------------------+--------------------+
|user_id|              review|              tokens|       refind_tokens|            features|         idfFeatures|
+-------+--------------------+--------------------+--------------------+--------------------+--------------------+
|      1|T really liked th...|[t, really, liked...|[really, liked, m...|(11,[0,1,4],[1.0,...|(11,[0,1,4],[0.0,...|
|      2|I would recommend...|[i, would, recomm...|[recommend, movie...|(11,[0,3,8],[1.0,...|(11,[0,3,8],[0.0,...|
|      3|movie was alright...|[movie, was, alri...|[movie, alright, ...|(11,[0,2,5,10],[1...|(11,[0,2,5,10],[0...|
|      4|I am never watchi...|[i, am, never, wa...|[never, watching,...|(11,[0,1,6,7,9],[...|(11,[0,1,6,7,9],[...|
+-------+--------------------+--------------------+--------------------+--------------------+--------------------+

['movie', 'liked', 'horrible', 'recommend', 'really', 'alright', 'never', 'watching', 'friends', 'ever', 'acting']
[0.         0.51082562 0.91629073 0.91629073 0.91629073 0.91629073
 0.91629073 0.91629073 0.91629073 0.91629073 0.91629073]

8:利用索引和IDF 值进行排序,并且将其显示为keyword得索引与IDF值对应

def func(partition):
    TOPK = 20
    for row in partition:
        # 利用索引和IDf值进行排序
        _ = list(zip( row.idfFeatures.indices,row.idfFeatures.values))
        _ = sorted(_,key=lambda x:x[1],reverse=True)
        result = _[:TOPK]
        for word_index,word_value in result:
            yield row.user_id,int(word_index),round(float(word_value),4)

_keywordsByTFIDF  = idf_df.rdd.mapPartitions(func).toDF(["user_id","index","tfidf"])
_keywordsByTFIDF .show()
+-------+-----+------+
|user_id|index| tfidf|
+-------+-----+------+
|      1|    4|0.9163|
|      1|    1|0.5108|
|      1|    0|   0.0|
|      2|    3|0.9163|
|      2|    8|0.9163|
|      2|    0|   0.0|
|      3|    2|0.9163|
|      3|    5|0.9163|
|      3|   10|0.9163|
|      3|    0|   0.0|
|      4|    6|0.9163|
|      4|    7|0.9163|
|      4|    9|0.9163|
|      4|    1|0.5108|
|      4|    0|   0.0|
+-------+-----+------+

9:寻找 keyword 和IDF 得对应关系
#cv_model.vocabulary, idf_model.idf.toArray()  分别存储 keyword 与其对应得idf值

keywords_list_with_idf = list(zip(cv_model.vocabulary, idf_model.idf.toArray())) print(keywords_list_with_idf)
[('movie', 0.0), ('liked', 0.5108256237659907), ('horrible', 0.9162907318741551), ('recommend', 0.9162907318741551), ('really', 0.9162907318741551), ('alright', 0.9162907318741551), ('never', 0.9162907318741551), ('watching', 0.9162907318741551), ('friends', 0.9162907318741551), ('ever', 0.9162907318741551), ('acting', 0.9162907318741551)]
datas = []
for i in range(len( keywords_list_with_idf)):
    
    word = keywords_list_with_idf[i][0]
    idf = float(keywords_list_with_idf[1][1])
    datas.append([word,idf,i])
    
print(datas)
[['movie', 0.5108256237659907, 0], ['liked', 0.5108256237659907, 1], ['horrible', 0.5108256237659907, 2], ['recommend', 0.5108256237659907, 3], ['really', 0.5108256237659907, 4], ['alright', 0.5108256237659907, 5], ['never', 0.5108256237659907, 6], ['watching', 0.5108256237659907, 7], ['friends', 0.5108256237659907, 8], ['ever', 0.5108256237659907, 9], ['acting', 0.5108256237659907, 10]]

将数组转换为RDD,再转换为DataFrame

oa = TestSparkCountVectorizer()
sc = oa.spark.sparkContext
rdd = sc.parallelize(datas)
df = rdd.toDF(["keywords", "idf", "index"])
df.show()
+---------+------------------+-----+
| keywords|               idf|index|
+---------+------------------+-----+
|    movie|0.5108256237659907|    0|
|    liked|0.5108256237659907|    1|
| horrible|0.5108256237659907|    2|
|recommend|0.5108256237659907|    3|
|   really|0.5108256237659907|    4|
|  alright|0.5108256237659907|    5|
|    never|0.5108256237659907|    6|
| watching|0.5108256237659907|    7|
|  friends|0.5108256237659907|    8|
|     ever|0.5108256237659907|    9|
|   acting|0.5108256237659907|   10|
+---------+------------------+-----+

 将其结果进行合并

df.join(_keywordsByTFIDF,df.index==_keywordsByTFIDF.index).select(["user_id","keywords","tfidf"]).show()
+-------+---------+------+
|user_id| keywords| tfidf|
+-------+---------+------+
|      1|    movie|   0.0|
|      2|    movie|   0.0|
|      3|    movie|   0.0|
|      4|    movie|   0.0|
|      1|    liked|0.5108|
|      4|    liked|0.5108|
|      3| horrible|0.9163|
|      2|recommend|0.9163|
|      1|   really|0.9163|
|      3|  alright|0.9163|
|      4|    never|0.9163|
|      4| watching|0.9163|
|      2|  friends|0.9163|
|      4|     ever|0.9163|
|      3|   acting|0.9163|
+-------+---------+------+

10:保存至数据库

df.write.insertInto("表名")