Python 2.7多线程优化:如何规避最慢任务拖慢整体执行速度?
你的问题核心在于当前用join()等待所有线程完成,导致总耗时被最慢的任务(10秒sleep)绑定。如果你的业务场景允许不等待所有任务结束就处理已完成的结果,或者想更高效地管理任务,可以用以下几种方案:
方案1:用Queue异步收集任务结果,边执行边处理
通过Queue让线程完成后把结果放入队列,主线程不用等待所有线程结束,而是持续从队列中取出已完成的任务结果进行处理,这样整体流程不会被慢任务阻塞。
代码示例:
import time from threading import Thread import random from Queue import Queue def task(counter, num_things, result_queue): item_num = counter + 1 sleep_for = time_sleep.get(item_num, {}).get("sleep", 1) print('Start task %s, sleeping for %s' % (item_num, sleep_for)) time.sleep(sleep_for) task_cost = random.randrange(50, 200) # 将结果放入队列 result_queue.put((item_num, task_cost)) print('Task %s completed, cost: %s' % (item_num, task_cost)) # 配置参数 num_things = 10 time_sleep = {1: {"sleep": 7}, 10: {"sleep": 10}} version = "thread-with-queue" num_main = 2 for _ in range(num_main): start_time = time.time() result_queue = Queue() jobs = [] # 启动所有线程 for things_counter in range(num_things): job = Thread(target=task, args=(things_counter, num_things, result_queue)) jobs.append(job) job.start() # 处理已完成的任务结果,不用等所有线程结束 completed = 0 while completed < num_things: item_num, cost = result_queue.get() # 这里可以添加对结果的处理逻辑,比如写入文件、统计等 print('Processing result: Task %s cost %s' % (item_num, cost)) completed += 1 end_time = time.time() print('Version: %s. Total time: %.2f seconds.' % (version, end_time - start_time))
这个方案中,主线程会在任务完成一个就处理一个,不用等最慢的任务结束,整体流程的“有效处理进度”不会被慢任务拖慢,总耗时还是由最慢任务决定,但中间可以及时处理结果。
方案2:用multiprocessing进程池,无序获取结果
如果你的任务是CPU密集型(当前IO密集型场景也适用),可以用multiprocessing.Pool的imap_unordered方法,它会优先返回已完成的任务结果,不需要按任务顺序等待。
代码示例:
import time import random from multiprocessing import Pool def task(counter): item_num = counter + 1 sleep_for = time_sleep.get(item_num, {}).get("sleep", 1) print('Start task %s, sleeping for %s' % (item_num, sleep_for)) time.sleep(sleep_for) task_cost = random.randrange(50, 200) print('Task %s completed, cost: %s' % (item_num, task_cost)) return (item_num, task_cost) # 配置参数 num_things = 10 time_sleep = {1: {"sleep": 7}, 10: {"sleep": 10}} version = "multiprocessing-pool" num_main = 2 for _ in range(num_main): start_time = time.time() # 创建进程池,大小可以根据CPU核心数设置 pool = Pool(processes=4) # 用imap_unordered无序获取结果 for result in pool.imap_unordered(task, range(num_things)): item_num, cost = result print('Processing result: Task %s cost %s' % (item_num, cost)) pool.close() pool.join() end_time = time.time() print('Version: %s. Total time: %.2f seconds.' % (version, end_time - start_time))
imap_unordered会在任务完成时立即返回结果,不用按提交顺序等待,这样你可以更快地处理已完成的任务,总耗时依然由最慢任务决定,但中间的结果处理不会被阻塞。
方案3:用gevent协程(第三方库)
对于IO密集型任务,协程的切换开销比线程更小,效率更高。需要先在Python2.7中安装gevent:pip install gevent
代码示例:
import time import random from gevent import spawn, joinall from gevent.monkey import patch_all # 打补丁,让标准库支持协程 patch_all() def task(counter): item_num = counter + 1 sleep_for = time_sleep.get(item_num, {}).get("sleep", 1) print('Start task %s, sleeping for %s' % (item_num, sleep_for)) time.sleep(sleep_for) task_cost = random.randrange(50, 200) print('Task %s completed, cost: %s' % (item_num, task_cost)) return (item_num, task_cost) # 配置参数 num_things = 10 time_sleep = {1: {"sleep": 7}, 10: {"sleep": 10}} version = "gevent-coroutine" num_main = 2 for _ in range(num_main): start_time = time.time() # 创建协程任务 tasks = [spawn(task, counter) for counter in range(num_things)] # 等待所有协程完成,也可以逐个检查状态处理结果 joinall(tasks) # 处理结果 for task_obj in tasks: item_num, cost = task_obj.value print('Processing result: Task %s cost %s' % (item_num, cost)) end_time = time.time() print('Version: %s. Total time: %.2f seconds.' % (version, end_time - start_time))
协程在IO等待时会自动切换到其他任务,所以整体的资源利用率更高,总耗时同样由最慢任务决定,但执行效率比线程更高,尤其是任务数量多的时候。
关键说明
如果你的业务必须等待所有任务完成才能进入下一步,那总耗时必然由最慢的任务决定——因为所有任务都需要完成。但上述方案可以让你在等待期间先处理已完成的任务结果,避免整体流程完全被慢任务卡住,提升资源利用率和响应速度。
内容的提问来源于stack exchange,提问作者mozman2

