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

Python多线程运行时动态分配CPU:大词表任务闲置CPU利用问询

解决方案

核心思路是把单个大词表拆分成多个可并行处理的子任务,而非一个文件对应一个进程。这样小文件处理完成后,闲置的CPU可以直接接手大文件的分片任务,避免资源浪费。

具体修改步骤

1. 新增文件分片函数

先实现一个函数,根据文件行数将大文件拆分为固定大小的分片,小文件直接作为单个任务:

def split_file_into_chunks(file_path, chunk_size=100000):
    chunks = []
    # 统计文件总行数
    with open(file_path, 'r', encoding='utf-8') as f:
        total_lines = sum(1 for _ in f)
    
    # 小文件直接作为单个任务
    if total_lines <= chunk_size:
        chunks.append((file_path, 0, total_lines))
    else:
        # 计算分片数量并拆分大文件
        num_chunks = (total_lines + chunk_size - 1) // chunk_size
        for i in range(num_chunks):
            start_line = i * chunk_size
            end_line = min((i+1)*chunk_size, total_lines)
            chunks.append((file_path, start_line, end_line))
    return chunks

2. 修改处理函数适配分片

将原来的start函数调整为处理指定行范围的分片任务:

def process_chunk(args):
    file_path, start_line, end_line = args
    processed_data = []
    with open(file_path, 'r', encoding='utf-8') as f:
        # 跳过前面的行,定位到分片起始位置
        for _ in range(start_line):
            next(f)
        # 处理分片内的每一行
        for _ in range(start_line, end_line):
            line = f.readline().strip()
            # 替换为你原来start函数中的处理逻辑
            # 例如词频统计、数据清洗等操作
            processed_data.append(line)
    return processed_data

3. 重构主逻辑,构建统一任务队列

将所有文件(无论大小)转换为分片任务,再用进程池执行:

import os
import multiprocessing

def main():
    PATH_TO_WORDLISTS = "./your_wordlist_directory"
    # 获取所有词表文件的完整路径,过滤非txt文件
    wordlist_files = [
        os.path.join(PATH_TO_WORDLISTS, fname) 
        for fname in os.listdir(PATH_TO_WORDLISTS)
        if fname.endswith('.txt')
    ]
    
    # 构建所有任务:小文件直接加入,大文件拆成分片后加入
    all_tasks = []
    CHUNK_SIZE = 100000  # 可根据服务器性能调整,比如10万行/分片
    for file in wordlist_files:
        all_tasks.extend(split_file_into_chunks(file, CHUNK_SIZE))
    
    # 启动进程池执行任务
    cpu_count = multiprocessing.cpu_count()
    with multiprocessing.Pool(cpu_count) as pool:
        # 若不需要按顺序获取结果,用imap_unordered效率更高
        results = pool.map(process_chunk, all_tasks)
    
    # (可选)合并同一个文件的分片处理结果
    merged_results = {}
    for res, task_args in zip(results, all_tasks):
        file_path = task_args[0]
        if file_path not in merged_results:
            merged_results[file_path] = []
        merged_results[file_path].extend(res)

if __name__ == "__main__":
    main()

关键优化点说明

  • 分片大小调整:CHUNK_SIZE可根据处理逻辑和内存情况灵活调整——如果每行处理占用内存大,就调小分片;反之可调大,减少进程调度开销。
  • 无序结果处理:若无需保留任务执行顺序,使用pool.imap_unordered(process_chunk, all_tasks)代替map,能更快获取已完成的任务结果,避免无效等待。
  • 内存控制:统计超大文件(千万级行)的总行数时,可改用逐行流式计数,避免一次性加载整个文件到内存。

内容的提问来源于stack exchange,提问作者nostripe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 01:53:18