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("表名")