Python中如何等待多个线程中的任意一个完成?
Python中如何等待多个线程中的任意一个完成?
嘿,这个需求我之前写并发任务的时候也碰到过!其实Python标准库就有现成的方案,完全不用自己折腾计数信号量,还能轻松知道是哪个线程完成了任务,甚至方便你后续追加新线程,给你两种实用的方法:
方法一:用concurrent.futures.ThreadPoolExecutor + as_completed(最推荐)
这个是我平时用得最多的方式,代码简洁还直观,as_completed方法会帮你监听所有提交的任务,只要有一个任务完成就立刻返回对应的Future对象,你能直接拿到任务结果,还能关联到具体的任务标识。
举个实际的例子:
from concurrent.futures import ThreadPoolExecutor, as_completed import time import random # 定义你的任务函数,这里模拟不同耗时的工作单元 def do_work(task_id): sleep_time = random.randint(1, 5) time.sleep(sleep_time) return f"任务{task_id}完成,耗时{sleep_time}秒" def main(): # 初始化线程池,比如先开3个线程 with ThreadPoolExecutor(max_workers=3) as executor: # 提交一批初始任务,把task_id和Future对象关联起来 task_futures = {executor.submit(do_work, i): i for i in range(1, 4)} # 循环处理完成的任务 for future in as_completed(task_futures): task_id = task_futures[future] try: result = future.result() print(result) # 这里可以处理完成的工作单元,比如根据结果决定要不要加新任务 # 比如我们追加一个新任务试试 new_task_id = max(task_futures.values()) + 1 new_future = executor.submit(do_work, new_task_id) task_futures[new_future] = new_task_id print(f"已追加新任务{new_task_id}") except Exception as e: print(f"任务{task_id}执行出错: {e}") # 如果不需要再追加任务,也可以在这里做终止判断 # 比如当任务总数达到10就停止 if len(task_futures) >= 10: break if __name__ == "__main__": main()
你看,这个例子里,主线程会一直等着,哪个任务先完成就先处理哪个,还能随时提交新任务到线程池里,完全符合你的需求。而且通过task_futures这个字典,你能轻松对应到完成的是哪个任务ID,排查问题也方便。
方法二:用原生threading + queue.Queue(手动控制更灵活)
如果你不想用concurrent.futures,用原生的线程模块也能实现。核心思路是让每个线程完成任务后,把自己的信息(比如任务ID、结果)放到一个队列里,主线程只需要阻塞在队列的get()方法上,一旦有线程把数据放进来,主线程就会被唤醒,拿到完成的线程信息。
例子如下:
import threading import queue import time import random def do_work(task_id, result_queue): sleep_time = random.randint(1, 5) time.sleep(sleep_time) # 任务完成后把结果和任务ID放进队列 result_queue.put(("success", task_id, f"耗时{sleep_time}秒")) def main(): result_queue = queue.Queue() threads = [] # 启动初始3个线程 for task_id in range(1, 4): t = threading.Thread(target=do_work, args=(task_id, result_queue)) t.start() threads.append((t, task_id)) processed_tasks = 0 # 循环处理完成的任务 while processed_tasks < 5: # 比如处理5个任务就停 # 阻塞等待,直到有线程把结果放进来 status, task_id, msg = result_queue.get() if status == "success": print(f"任务{task_id}完成,{msg}") # 可以在这里启动新线程 new_task_id = max([tid for _, tid in threads]) + 1 new_t = threading.Thread(target=do_work, args=(new_task_id, result_queue)) new_t.start() threads.append((new_t, new_task_id)) print(f"已启动新任务{new_task_id}") processed_tasks += 1 # 最后等待所有线程完成(可选) for t, _ in threads: t.join() if __name__ == "__main__": main()
这种方法的好处是你能完全掌控线程的创建和结果传递的细节,适合需要自定义线程行为的场景。
两种方法都能完美解决“等待任意一个线程完成”的需求,还能满足你后续追加新线程的要求,比自己写信号量靠谱多啦~
备注:内容来源于stack exchange,提问作者Edward Falk
相关产品推荐
相关产品推荐

