Call back DP05 处理新字段 - kWords
多进程程序:多进程无法在Jupyter Notebook中运行,所以这里存放PyCharm的运行代码。
本程序注意点:
待处理程序具有返回值怎么办? :执行多进程的代码块前面,定义的变量
res是用于处理带有返回值return语句的变量;返回值无法pickle到本地怎么办? :
pickle序列化的对象具有一定的要求,本程序直接pickle时报错,无法pickle"lock"对象,所以这里我们多了一步赋值语句,解决了这个问题。
运行代码
import pickle
# Input: 变量名,保存路径
# Output: 保存路径
def save_variable1(v, filename):
f = open(filename, 'wb')
pickle.dump(v, f, 0) # 在dump()里加上第三个参数,设为0(ASCII协议).
f.close()
return filename
# Input: 保存的变量路径
# Output: 变量
def load_variable1(filename):
f = open(filename, 'rb')
r = pickle.load(f)
f.close()
return r
def save_variable2(v, filename):
try:
# f = open('test001.txt','w',encoding='utf-8') #open 打开要写入的文件,'w'就是写入的意思,如果有中文encoding进行编码
with open(filename,'w',encoding='utf-8') as f:
f.write(v) # write是写入的意思,将()中的内容写入到f文件里
# 无论try中的程序是否存在报错,则都执行下面的关闭
finally:
if f: # 判断下f是否打开了,如果没打开则不需要关闭,打开了则关闭
f.close()
return filename
def load_variable2(filename):
# 首先try一下,如果程序打开以及在打开后读取一系列操作后有报错,则不中断程序
try:
# f = open(filename, 'r', encoding='utf-8')
with open(filename, 'r', encoding='utf-8') as f:
strContent = f.read()
# 无论try中的程序是否存在报错,则都执行下面的关闭
finally:
if f: # 判断下f是否打开了,如果没打开则不需要关闭,打开了则关闭
f.close()
return strContent
import pymongo
import time
import os # ▲
from multiprocessing import Pool
import re
# ▲ 加了一个进程的名字
def computeWordFreq_pro(curCol, targetDB, targetCol, per = 1.0):
loc = locals()
db = pymongo.MongoClient("localhost:27017", serverSelectionTimeoutMS=10000)
exec(targetDB + " = db['" + targetDB + "']")
print(targetDB + ": 已连接!")
print(f"我是{curCol}子进程,我的进程号为{os.getpid()},我开始执行了...") # ▲
exec(targetCol + " = " + targetDB + "['" + targetCol + "']")
exec("recs = " + targetCol + ".find()")
tempCnt_Words = 0 # ●1 存放当前集合中的词的总数
freqWords = {} # key(词汇),value(词频数)
for cW in loc['recs']:
# 手动写词频统计函数,基于字典实现
currentContent = cW.get('kWords')
tempCnt_Words = tempCnt_Words + len(currentContent) # 存放当前集合中关键词的总数量
for eachWord in currentContent:
freqWords[eachWord] = freqWords.get(eachWord, 0) + 1 # 统计每个关键词的词频
items = list(freqWords.items()) # 转换为list
total_num = len(items) # ●2 不一样的词的个数
# print("items.len = " + str(total_num)) # ●2 不一样的词的个数
items.sort(key=lambda x: x[1], reverse=True) # 按照词频排序
topK = int(total_num * per)
wordcloudDic = dict() # ●3 存放前k个词 & 词频
for i in range(topK):
word, count = items[i]
wordcloudDic[word] = count
# print("{0:<20}{1:>10}".format(word, count))
res = list()
res.append(targetCol)
res.append(tempCnt_Words)
res.append(len(items))
res.append(wordcloudDic)
print(f"我是{curCol}子进程,我执行完了!") # ▲
return res
if __name__ == '__main__':
targetDB = 'db_Weibo_SH0331'
startTime = time.time()
print(time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(startTime)))
# 开始处理----↓
patternNum = '[0-9]+'
# 循环处理24个集合
subCol24 = list()
subCol24.append(["col_weibo_timePeriod_" + str(period) for period in range(0, 24)])
colPackages = subCol24[0]
per = 0.2
p = Pool(5)
res = list()
for i in range(0, len(colPackages)):
post = re.search(patternNum, colPackages[i]).group()
res.append(p.apply_async(computeWordFreq_pro, args=('col' + post, targetDB, colPackages[i], per)))
print('Waiting for all subprocesses done...')
p.close()
p.join()
'''
process_list = [] # 存放开启的进程
for i in range(0, len(colPackages)):
post = re.search(patternNum, colPackages[i]).group()
# 进程中的参数args表示调用对象的位置参数元组.注意:元组中只有一个元素时结尾要加","逗号
p = Process(target=filterSentences, args=('col' + post, targetDB, colPackages[i], fields))
p.start()
print('------> Thread id : %d' % p.ident) # 进程id
print('------> Thread name : %s' % p.name) # 进程name
process_list.append(p)
for i in process_list:
i.join() # 阻塞每个子进程,主进程会等待所有子进程结束再结束主进程
'''
print("主进程结束!")
# 结束处理----↑
endTime = time.time()
print(time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(endTime)))
print("Time consuming(s): %.2f" % (endTime - startTime))
wordFreqDict_MulPro = dict()
for r in res:
content = r.get()
# print(content)
# ['col_weibo_timePeriod_0', 1305921, 36127, {'疫情': 47415, '医院': 40077, '小区': 38743, ...}]
wordFreqDict_MulPro[content[0]] = content[1:]
print(type(wordFreqDict_MulPro))
save_variable1(wordFreqDict_MulPro, 'wordFreqDict_MulPro_WithType.txt')
save_variable2(str(wordFreqDict_MulPro), 'wordFreqDict_MulPro_WithoutType.txt') # 只能查看字符串
运行结果
2022-09-03 15:38:09
Waiting for all subprocesses done...
db_Weibo_SH0331: 已连接!
我是col2子进程,我的进程号为32420,我开始执行了...
db_Weibo_SH0331: 已连接!
我是col0子进程,我的进程号为28736,我开始执行了...
db_Weibo_SH0331: 已连接!
我是col4子进程,我的进程号为31112,我开始执行了...
db_Weibo_SH0331: 已连接!
我是col3子进程,我的进程号为10556,我开始执行了...
db_Weibo_SH0331: 已连接!
我是col1子进程,我的进程号为2524,我开始执行了...
我是col4子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col5子进程,我的进程号为31112,我开始执行了...
我是col3子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col6子进程,我的进程号为10556,我开始执行了...
我是col5子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col7子进程,我的进程号为31112,我开始执行了...
我是col2子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col8子进程,我的进程号为32420,我开始执行了...
我是col6子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col9子进程,我的进程号为10556,我开始执行了...
我是col1子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col10子进程,我的进程号为2524,我开始执行了...
我是col7子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col11子进程,我的进程号为31112,我开始执行了...
我是col0子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col12子进程,我的进程号为28736,我开始执行了...
我是col8子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col13子进程,我的进程号为32420,我开始执行了...
我是col9子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col14子进程,我的进程号为10556,我开始执行了...
我是col10子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col15子进程,我的进程号为2524,我开始执行了...
我是col12子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col16子进程,我的进程号为28736,我开始执行了...
我是col11子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col17子进程,我的进程号为31112,我开始执行了...
我是col14子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col18子进程,我的进程号为10556,我开始执行了...
我是col13子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col19子进程,我的进程号为32420,我开始执行了...
我是col15子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col20子进程,我的进程号为2524,我开始执行了...
我是col16子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col21子进程,我的进程号为28736,我开始执行了...
我是col17子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col22子进程,我的进程号为31112,我开始执行了...
我是col18子进程,我执行完了!
db_Weibo_SH0331: 已连接!
我是col23子进程,我的进程号为10556,我开始执行了...
我是col19子进程,我执行完了!
我是col20子进程,我执行完了!
我是col22子进程,我执行完了!
我是col21子进程,我执行完了!
我是col23子进程,我执行完了!
主进程结束!
2022-09-03 15:38:51
Time consuming(s): 42.13
<class 'dict'>
Process finished with exit code 0
对比发现,多进程速度有明显提升。
多进程耗时:
Time consuming(s): 42.13普通运行耗时:
Time consuming: 115.33
处理得到的部分结果(文档内容展示)
{'col_weibo_timePeriod_0': [1305921, 36127, {'疫情': 47415, '医院': 40077, '小区': 38743, '30': 37617, '病人': 37557, '2022': 37474, '120': 37039, '离世': 36322, '哮喘': 36315, '救治': 36302, '救助': 36275, '110': 36208, '救护车': 36186, '邻居': 36181, '发作': 36126, '阻拦': 36124, '804': 36121, '40': 36083, '202': 35334, '11': 25890, '上海': 12816, '核酸': 7762, '微博': 7687, '到达': 7041, '转发': 5379, '##': 5114, '新冠': 4188, '浦东': 4043, '封控': 3782, '筛查': 3561, '全域': 3119, '全员': 3073, '血透': 2903, '求助': 2696, '防疫': 2544, '静态': 2513, ..., '来得': 6}],
'col_weibo_timePeriod_1': [809348, 27751, {'疫情': 28547, '医院': 25508, '小区': 24719, '120': 24341, '病人': 23647, '30': 22599, '救助': 22472, '2022': 22432, '离世': 22317, '哮喘': 22290, '救治': 22192, '邻居': 22133, '救护车': 22119, '110': 22103, '阻拦': 22099, '发作': 22064, '804': 22058, '40': 21986, '202': 21506, '11': 16563, '上海': 10039, '核酸': 6037, ...}],
...,
'col_weibo_timePeriod_23': [744825, 36738, {'疫情': 16786, '上海': 11156, '微博': 7294, '核酸': 7070, '新冠': 7063, '小区': 7033, '120': 6626, '居委': 6300, '人员': 6154, '28': 5544, '我们': 5465, '救助': 5080, '不能': 5000, '抗疫': 4881, '30': 4849, '医生': 4824, '救护车': 4793, '大家': 4787, '16': 4776, '邻居': 4760, '看到': 4444, '父亲': 4418, 'AED': 4381, '救援': 4279, '202': 4221, '上车': 4208, '病毒': 4117, '奥密克戎': 3900, '检测': 3857, '110': 3756, 'XE': 3596, '阳性': 3587, '##': 3579, '31': 3550, '毒株': 3527, '转发': 3410, 'BA.2': 3409, '传播速度': 3376, '变体': 3365, '10%': 3365, '重组': 3360, '637': 3358, 'UKHSA': 3292, '封控': 3290, '医院': 3221, '奥密克': 3219, '戎子': 3162, 'BA.1': 3143, '血透': 3120, '问能': 3076, '9.8%': 2988, '小时': 2977, '防控': 2830, '两种': 2802, '2022': 2463, '求助': 2208, '浦东': 2142, '24': 2088, '隔离': 2086, '20': 2011, '求求': 1962, '防疫': 1915, ..., '字号': 6}]}