多进程多线程场景下全局列表数据无法跨进程传递问题
多进程+多线程架构下全局列表为空的问题解决
问题描述
我在Python中实现多进程结合多线程的架构,计划创建近10个线程,当前使用2个线程测试。每个线程内的计算结果会追加到全局列表中,在创建线程的run_process函数中可正常打印列表数据,但调用plot_graph函数时,这些全局列表却为空。尽管已将列表声明为全局变量,该问题仍存在。相关代码如下:
# importing package from datasets import load_dataset from transformers import pipeline import matplotlib.pyplot as plt from collections import Counter from nltk.translate import bleu from nltk import cluster import pandas as pd from threading import Thread from multiprocessing import Process from Levenshtein import ratio global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore global sum_model generatedSumamrycosineScore = [] generatedSumamryjaccordScore = [] originalSummaryLevenshteinScore = [] originalSummaryScore = [] sum_model = pipeline('summarization', model='facebook/bart-large-cnn') def findCosine_Similarity(article, summary): counter1 = Counter(article.split()) counter2 = Counter(summary.split()) all_items = set(counter1.keys()).union(set(counter2.keys())) sim = [] v1 = [counter1[k] for k in all_items] for b in summary.split('. '): if b.strip() != '': c2 = Counter(b.split()) v2 = [c2[k] for k in all_items] sim.append({'article_line': article, 'summary_line': b, "similarity": cluster.util.cosine_distance(v1, v2)}) df = pd.DataFrame(sim) df_sum = df[abs(df['similarity']) >= 0.75] s = '. '.join(df_sum['summary_line'].to_list()) return bleu([s.split()], article.split()) def jaccard_similarity(art, summary): s1 = set(art.split(' ')) sum = [] for b in summary.split('. '): if b.strip() != '': s2 = set(b.split(' ')) sum.append({'article_line': art, 'summary_line': b, "similarity": float( len(s1.intersection(s2)) / len(s1.union(s2)))}) df = pd.DataFrame(sum) df_sum = df[abs(df['similarity']) <= 0.25] s = '. '.join(df_sum['summary_line'].to_list()) return bleu([s.split()], art.split()) def Levenshtein_similarity(art, summary): sum = [] for b in summary.split('. '): if b.strip() != '': sum.append({'article_line': art, 'summary_line': b, "similarity": ratio(art, b)}) df = pd.DataFrame(sum) df_sum = df[abs(df['similarity']) >= 0.75] s = '. '.join(df_sum['summary_line'].to_list()) return bleu([s.split()], art.split()) def PredictSummary(*art): try: global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore global sum_model l1 = art[0] summary = sum_model(l1) l2 = summary[0]['summary_text'] generatedSumamrycosineScore.append(findCosine_Similarity(l1, l2)) generatedSumamryjaccordScore.append(jaccard_similarity(l1, l2)) originalSummaryLevenshteinScore.append(Levenshtein_similarity(l1, l2)) originalSummaryScore.append( bleu([art[1].split()], l1.split())) except: pass def run_process(*art): threadlist = [] print("creating Threads") threadlist.append(Thread(target=PredictSummary, args=list(art[0].values())[:2])) threadlist.append(Thread(target=PredictSummary, args=list(art[1].values())[:2])) for t in threadlist: t.start() for t in threadlist: t.join() print("Threads created") global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore # Printing values here print(generatedSumamrycosineScore, generatedSumamryjaccordScore, originalSummaryLevenshteinScore, originalSummaryScore) def plot_graph(): global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore # list values are empty here print(generatedSumamrycosineScore, generatedSumamryjaccordScore, originalSummaryLevenshteinScore, originalSummaryScore) plt.plot(generatedSumamryjaccordScore, label="generatedSumamryjaccordScore") plt.plot(generatedSumamrycosineScore, label="generatedSumamrycosineScore") plt.plot(originalSummaryLevenshteinScore, label="originalSummaryLevenshteinScore") plt.plot(originalSummaryScore, label="originalSummaryScore") plt.legend() plt.show() plt.savefig("Bleugraph.png", dpi=1000) if __name__ == '__main__': datasets = load_dataset('cnn_dailymail', '3.0.0') i = 0 while i < 1000: try: print("started", i) art = [] for k in range(0, 30): art.append(datasets['test'][i+k]) processlist = [] processlist.append(Process(target=run_process, args=art[:10])) for t in processlist: t.start() for t in processlist: t.join() break i += 30 except: continue plot_graph()
原因分析
问题核心在于多进程的内存隔离机制:Python的multiprocessing.Process创建的子进程会完整复制父进程的内存空间,包括全局变量的初始值。子进程中对全局列表的修改,都是在自己的内存副本上操作,父进程的原列表完全不受影响。run_process在子进程中执行,打印的是子进程里的列表;而plot_graph在父进程中调用,访问的是父进程初始的空列表,所以显示为空。
解决方案
使用multiprocessing.Manager创建可跨进程共享的列表,所有进程操作的是同一个列表实例,适配现有代码改动较小。
修改后的完整代码
# importing package from datasets import load_dataset from transformers import pipeline import matplotlib.pyplot as plt from collections import Counter from nltk.translate import bleu from nltk import cluster import pandas as pd from threading import Thread from multiprocessing import Process, Manager from Levenshtein import ratio # 用Manager创建跨进程共享的列表 manager = Manager() generatedSumamrycosineScore = manager.list() generatedSumamryjaccordScore = manager.list() originalSummaryLevenshteinScore = manager.list() originalSummaryScore = manager.list() sum_model = pipeline('summarization', model='facebook/bart-large-cnn') def findCosine_Similarity(article, summary): counter1 = Counter(article.split()) counter2 = Counter(summary.split()) all_items = set(counter1.keys()).union(set(counter2.keys())) sim = [] v1 = [counter1[k] for k in all_items] for b in summary.split('. '): if b.strip() != '': c2 = Counter(b.split()) v2 = [c2[k] for k in all_items] sim.append({'article_line': article, 'summary_line': b, "similarity": cluster.util.cosine_distance(v1, v2)}) df = pd.DataFrame(sim) df_sum = df[abs(df['similarity']) >= 0.75] s = '. '.join(df_sum['summary_line'].to_list()) return bleu([s.split()], article.split()) def jaccard_similarity(art, summary): s1 = set(art.split(' ')) sum = [] for b in summary.split('. '): if b.strip() != '': s2 = set(b.split(' ')) sum.append({'article_line': art, 'summary_line': b, "similarity": float( len(s1.intersection(s2)) / len(s1.union(s2)))}) df = pd.DataFrame(sum) df_sum = df[abs(df['similarity']) <= 0.25] s = '. '.join(df_sum['summary_line'].to_list()) return bleu([s.split()], art.split()) def Levenshtein_similarity(art, summary): sum = [] for b in summary.split('. '): if b.strip() != '': sum.append({'article_line': art, 'summary_line': b, "similarity": ratio(art, b)}) df = pd.DataFrame(sum) df_sum = df[abs(df['similarity']) >= 0.75] s = '. '.join(df_sum['summary_line'].to_list()) return bleu([s.split()], art.split()) def PredictSummary(*art): try: global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore global sum_model l1 = art[0] summary = sum_model(l1) l2 = summary[0]['summary_text'] generatedSumamrycosineScore.append(findCosine_Similarity(l1, l2)) generatedSumamryjaccordScore.append(jaccard_similarity(l1, l2)) originalSummaryLevenshteinScore.append(Levenshtein_similarity(l1, l2)) originalSummaryScore.append( bleu([art[1].split()], l1.split())) except: pass def run_process(*art): threadlist = [] print("creating Threads") threadlist.append(Thread(target=PredictSummary, args=list(art[0].values())[:2])) threadlist.append(Thread(target=PredictSummary, args=list(art[1].values())[:2])) for t in threadlist: t.start() for t in threadlist: t.join() print("Threads created") global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore # Printing values here print(generatedSumamrycosineScore, generatedSumamryjaccordScore, originalSummaryLevenshteinScore, originalSummaryScore) def plot_graph(): global generatedSumamrycosineScore global generatedSumamryjaccordScore global originalSummaryLevenshteinScore global originalSummaryScore # 转换为普通列表方便绘图(可选,matplotlib也支持manager.list) cosine_scores = list(generatedSumamrycosineScore) jaccard_scores = list(generatedSumamryjaccordScore) levenshtein_scores = list(originalSummaryLevenshteinScore) original_scores = list(originalSummaryScore) print(cosine_scores, jaccard_scores, levenshtein_scores, original_scores) plt.plot(jaccard_scores, label="generatedSumamryjaccordScore") plt.plot(cosine_scores, label="generatedSumamrycosineScore") plt.plot(levenshtein_scores, label="originalSummaryLevenshteinScore") plt.plot(original_scores, label="originalSummaryScore") plt.legend() plt.show() plt.savefig("Bleugraph.png", dpi=1000) if __name__ == '__main__': datasets = load_dataset('cnn_dailymail', '3.0.0') i = 0 while i < 1000: try: print("started", i) art = [] for k in range(0, 30): art.append(datasets['test'][i+k]) processlist = [] processlist.append(Process(target=run_process, args=art[:10])) for t in processlist: t.start() for t in processlist: t.join() break i += 30 except: continue plot_graph()
关键改动说明
- 导入
multiprocessing.Manager,用manager.list()创建四个共享列表,替代原来的普通全局列表 - 子进程中对共享列表的修改会直接作用于父进程中的同一个列表实例
plot_graph中可以直接访问共享列表,也可以转换为普通列表(避免潜在的多进程同步问题)
内容的提问来源于stack exchange,提问作者Hammad Asif
相关产品推荐
相关产品推荐

