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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 12:00:27