Python实现多线程/多进程并发执行避免串行延迟
问题原因
- 原代码调用
starmap_async(...).get()会直接阻塞主线程,必须等所有profile对应的任务全部执行完,才会进入后面的for循环;而你把取Last_name、while循环的核心逻辑全写在了主线程的串行for循环里,根本没有放到并行的工作单元中,自然是逐个顺序执行。 - 原代码存在语法错误:字典取值
[Last_name]没有给Last_name加引号,运行会直接报NameError;修改后的多进程版本把工作函数嵌套定义在main函数里,Windows环境下会因为函数无法序列化直接报错。 - 虽然线程池/进程池开了和profile数量一致的工作数,但核心业务逻辑没放到工作函数里,并行资源没有真正被利用。
最小延迟并行实现方案
要做到所有profile任务同步启动、执行延迟最小,核心是把每个profile对应的完整业务逻辑全部放到工作函数中,提交任务后不等待全部任务完成,谁先执行完就先处理谁的结果,避免无意义的阻塞等待。
多线程版本(优先选,适合IO密集场景:网络请求、文件读写等绝大多数业务场景,线程开销更小)
import json import threading import time import random from collections import OrderedDict from multiprocessing.pool import ThreadPool def process_profile(profile_key, profile_data): # 模拟业务处理延迟 time.sleep(random.random()) last_name = profile_data["Last_name"] print(f"线程[{threading.current_thread().name}] 处理{profile_key}完成, Last_name: {last_name}") # 原代码中的while循环逻辑直接放在工作线程内执行,不需要等所有任务跑完再串行处理 timeout = time.time() while True: # 替换成实际业务逻辑 print(f"{profile_key} 运行中, Last_name: {last_name}") time.sleep(1) if time.time() - timeout > 60: break return {profile_key: last_name} if __name__ == '__main__': with open('profileMulti.json', 'r', encoding='UTF-8') as json_data: config = json.load(json_data, object_pairs_hook=OrderedDict) profiles = list(config.get('Profiles', {}).items()) worker_count = len(profiles) # 线程数和profile数一致,保证所有任务提交后立刻被拾取执行,无排队 pool = ThreadPool(worker_count) print("所有profile任务同步启动") # imap_unordered会优先返回先完成的任务结果,端到端延迟最低 for result in pool.imap_unordered(lambda item: process_profile(item[0], item[1]), profiles): print(f"收到任务返回结果: {result}") pool.close() pool.join()
多进程版本(适合CPU密集场景:大量数值计算、数据解析等,可绕过GIL利用多核并行)
import json import os import time import random from collections import OrderedDict from multiprocessing import Pool # 多进程工作函数必须定义在顶层,不能嵌套,否则无法序列化 def process_profile(profile_item): profile_key, profile_data = profile_item # 模拟业务处理延迟 time.sleep(random.random()) last_name = profile_data["Last_name"] print(f"进程[{os.getpid()}] 处理{profile_key}完成, Last_name: {last_name}") # 原代码中的while循环逻辑直接放在工作进程内执行 timeout = time.time() while True: # 替换成实际业务逻辑 print(f"{profile_key} 运行中, Last_name: {last_name}") time.sleep(1) if time.time() - timeout > 60: break return {profile_key: last_name} if __name__ == '__main__': with open('profileMulti.json', 'r', encoding='UTF-8') as json_data: config = json.load(json_data, object_pairs_hook=OrderedDict) profiles = list(config.get('Profiles', {}).items()) worker_count = len(profiles) pool = Pool(worker_count) print("所有profile任务同步启动") for result in pool.imap_unordered(process_profile, profiles): print(f"收到任务返回结果: {result}") pool.close() pool.join()
关键改动说明
- 所有和单个profile绑定的业务逻辑全部移入工作函数,任务提交后立刻被工作线程/进程拾取,所有任务几乎同时启动,不存在串行等待。
- 移除了
starmap_async后的阻塞get()调用,改用imap_unordered获取结果,任务完成一个就返回一个,没有批量等待的开销,延迟最小。 - 修复了原代码的语法错误、多进程函数嵌套的兼容性问题,符合Python并行编程规范。
内容的提问来源于stack exchange,提问作者turtlecoder
相关产品推荐
相关产品推荐

