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
相关产品推荐
相关产品推荐

