You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为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   |
+----+--------------+---------------+----------------+-------------+----------------+

核心问题分析

  1. GIL限制:Python全局解释器锁(GIL)导致CPU密集型任务无法通过多线程实现真正并行,线程切换反而会增加额外开销,这是多线程未提速的核心原因。
  2. 线程池重复创建:每次处理一个chunk就新建ThreadPoolExecutor,频繁的线程创建销毁会放大开销。
  3. 数据拆分低效:chunk[i::num_threads]的拆分方式让线程处理分散的行,缓存命中率低,影响计算效率。
  4. 小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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 08:47:04