为Python函数添加多线程后运行时长未缩短,求优化方案
Python多线程词频统计优化求助
我正在开发一个文本词频统计的Python程序,目标是按线程数拆分数据块实现负载均衡。刚接触Python多线程,不确定当前实现是否最优。加了多线程后运行时长没缩短,附上代码和测试输出,求优化建议。
代码
#so far this is the most optimal import heapq import re import time import psutil from collections import defaultdict import requests import tracemalloc from concurrent.futures import ThreadPoolExecutor import os # Load stop words stop_words_url = "https://gist.githubusercontent.com/sebleier/554280/raw/7e0e4a1ce04c2bb7bd41089c9821dbcf6d0c786c/NLTK's%2520list%2520of%2520english%2520stopwords" stop_words = set(requests.get(stop_words_url).text.split()) def tokenize(line): return re.findall(r'\b\w+\b', line) def count_words(chunk): word_count = defaultdict(int) for line in chunk: words = tokenize(line) for word in words: if word.lower() not in stop_words: word_count[word] += 1 return word_count def top_k_words(word_count, k): return heapq.nlargest(k, word_count.items(), key=lambda x: x[1]) def analyze_performance(file_path, k=10, chunk_size=10000, num_threads=os.cpu_count()): start_time = time.time() tracemalloc.start() word_count = defaultdict(int) with open(file_path, 'r', encoding='utf-8') as file: chunk = [] for line in file: chunk.append(line) if len(chunk) == chunk_size: with ThreadPoolExecutor(max_workers=num_threads) as executor: futures = {executor.submit(count_words, chunk[i::num_threads]) for i in range(num_threads)} for future in futures: chunk_word_count = future.result() for word, count in chunk_word_count.items(): word_count[word] += count chunk = [] if chunk: with ThreadPoolExecutor(max_workers=num_threads) as executor: futures = {executor.submit(count_words, chunk[i::num_threads]) for i in range(num_threads)} for future in futures: chunk_word_count = future.result() for word, count in chunk_word_count.items(): word_count[word] += count top_k = top_k_words(word_count, k) current, peak = tracemalloc.get_traced_memory() tracemalloc.stop() end_time = time.time() elapsed_time = end_time - start_time cpu_percent = psutil.cpu_percent() print(f"Top {k} words: {top_k}") print(f"Elapsed time: {elapsed_time:.2f} seconds") print(f"CPU usage: {cpu_percent}%") #print(f"Memory usage: {memory_usage / (1024 * 1024):.2f} MB") return elapsed_time, cpu_percent, peak if __name__ == "__main__": file_paths = [ "small_50MB_dataset.txt", ] chunk_sizes = [10, 100, 1000, 10000, 100000, 1000000, 10000000, 100000000, 1000000000, 10000000000] results = [] for file_path in file_paths: print(f"Processing {file_path}") for chunk_size in chunk_sizes: for numThreads in range(1, os.cpu_count() + 1): print("Partition Size:", chunk_size / (1024 * 1024), "MB", "chunk size of:", chunk_size) elapsed_time, cpu_usage, memory_usage = analyze_performance(file_path, chunk_size=chunk_size, num_threads=numThreads) print("\n") result = { "chunk_size": chunk_size, "num_threads": numThreads, "elapsed_time": elapsed_time, "cpu_usage": cpu_usage, "memory_usage": float(memory_usage / 10**6) } results.append(result) # Create pandas DataFrame from the results import pandas as pd df2 = pd.DataFrame(results)
测试输出
+----+--------------+---------------+----------------+-------------+----------------+ | | chunk_size | num_threads | elapsed_time | cpu_usage | memory_usage | |----+--------------+---------------+----------------+-------------+----------------| | 40 | 10 | 1 | 51.2827 | 24.8 | 8.01863 | | 41 | 10 | 2 | 60.1906 | 65.5 | 8.3454 | | 42 | 100 | 1 | 32.4096 | 64.4 | 8.11009 | | 43 | 100 | 2 | 33.402 | 60 | 8.16907 | | 44 | 1000 | 1 | 25.7621 | 62.5 | 8.48084 | | 45 | 1000 | 2 | 31.2087 | 65 | 9.02304 | | 46 | 10000 | 1 | 24.5674 | 70.6 | 12.702 | | 47 | 10000 | 2 | 23.1408 | 63.7 | 13.9474 | | 48 | 100000 | 1 | 19.4707 | 58.7 | 43.1203 | | 49 | 100000 | 2 | 21.5641 | 64.6 | 42.0958 | | 50 | 1e+06 | 1 | 21.23 | 61.9 | 99.1393 | | 51 | 1e+06 | 2 | 21.2195 | 60.7 | 104.215 | | 52 | 1e+07 | 1 | 21.5565 | 64.3 | 99.153 | | 53 | 1e+07 | 2 | 22.712 | 66.1 | 104.216 | | 54 | 1e+08 | 1 | 20.8239 | 61.9 | 99.1389 | | 55 | 1e+08 | 2 | 22.5298 | 63.9 | 104.217 | | 56 | 1e+09 | 1 | 21.4913 | 64.3 | 99.1535 | | 57 | 1e+09 | 2 | 20.9633 | 58.6 | 104.232 | | 58 | 1e+10 | 1 | 21.4864 | 64.6 | 99.1389 | | 59 | 1e+10 | 2 | 22.0327 | 63.9 | 104.216 | +----+--------------+---------------+----------------+-------------+----------------+
核心问题分析
- GIL限制:Python全局解释器锁(GIL)导致CPU密集型任务无法通过多线程实现真正并行,线程切换反而会增加额外开销,这是多线程未提速的核心原因。
- 线程池重复创建:每次处理一个chunk就新建
ThreadPoolExecutor,频繁的线程创建销毁会放大开销。 - 数据拆分低效:
chunk[i::num_threads]的拆分方式让线程处理分散的行,缓存命中率低,影响计算效率。 - 小Chunk开销过高:当chunk_size极小时(如10、100),线程调度和结果合并的开销远大于并行收益,导致多线程比单线程更慢。
具体优化方案
1. 改用多进程替代多线程
对于CPU密集型的词频统计,用ProcessPoolExecutor绕过GIL,实现真正的并行计算。
2. 复用线程/进程池
不要每次处理chunk都创建新池,在函数初始化时创建一次,重复使用。
3. 优化数据拆分逻辑
将chunk拆分为连续的子块分配给不同进程/线程,比如chunk[i*len(chunk)//num_workers : (i+1)*len(chunk)//num_workers],提升缓存利用率。
4. 合理设置Chunk大小
从测试结果看,chunk_size在10000左右时效果相对最优,过小会放大调度开销,过大则会占用过多内存,建议根据文件大小选择10000-100000行的chunk。
5. 优化停用词判断逻辑
将停用词提前转换为全小写,避免每次判断时重复执行word.lower()。
6. 简化结果合并
用collections.Counter的update方法替代手动循环累加,代码更简洁且效率更高。
优化后代码片段
from concurrent.futures import ProcessPoolExecutor from collections import Counter def analyze_performance(file_path, k=10, chunk_size=10000, num_workers=os.cpu_count()): start_time = time.time() tracemalloc.start() total_count = Counter() # 只创建一次进程池 with ProcessPoolExecutor(max_workers=num_workers) as executor: with open(file_path, 'r', encoding='utf-8') as file: chunk = [] for line in file: chunk.append(line) if len(chunk) == chunk_size: # 连续拆分chunk splits = [chunk[i*len(chunk)//num_workers : (i+1)*len(chunk)//num_workers] for i in range(num_workers)] # 批量提交任务并合并结果 for result in executor.map(count_words, splits): total_count.update(result) chunk = [] # 处理剩余数据 if chunk: splits = [chunk[i*len(chunk)//num_workers : (i+1)*len(chunk)//num_workers] for i in range(num_workers)] for result in executor.map(count_words, splits): total_count.update(result) top_k = top_k_words(total_count, k) # 后续性能统计代码保持不变...
额外建议
- 避免设置过大的chunk_size(如1e9、1e10),会导致一次性加载过多数据到内存,反而降低效率。
- 使用
psutil.cpu_percent(interval=1)代替无参数调用,能更准确反映程序运行期间的CPU使用率。
内容的提问来源于stack exchange,提问作者Adam Graham
相关产品推荐
相关产品推荐

