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

如何基于自定义线程与对象实现ThreadPoolExecutor风格并发控制

Python 原生节点任务调度实现方案

核心思路

要保证每个节点同一时间仅执行一条任务、空闲时自动领取新任务,不需要复杂的锁逻辑:

  • 用标准库queue.Queue作为线程安全的公共任务池,所有待执行任务统一入队
  • 每个节点绑定唯一的专属工作线程,线程循环从公共队列领取任务执行,因为单线程串行处理该节点的所有任务,天然避免同节点并发执行任务的冲突
  • 利用队列自带的阻塞等待机制,所有任务执行完成后自动退出调度

完整实现代码

from time import sleep
from threading import Thread
from queue import Queue


class Mork():
    def __init__(self, name) -> None:
        self.name = name
        self.introduction = f"Hi I'm {self.name}"
    
    def speak(self):
        print(f"{self.introduction}.")
        sleep(1)
    
    def scream(self):
        print(f"{self.introduction.upper()}!")
        sleep(2)
    
    def whisper(self):
        print(f"{self.introduction.lower()}...")
        sleep(3)


class Processor(object):
  """Process the data."""

  def __init__(self, node):
    self.node = node

  def __call__(self, mode):
    """Do something"""
    if mode == "speak":
        self.node.speak()
    elif mode == "scream":
        self.node.scream()
    elif mode == "whisper":
        self.node.whisper()
    else:
        print(f"{self.node.name} doesn't know how to {mode}...")


def complete_all_tasks_using_threads(processor_list, tasks):
    task_queue = Queue()
    # 所有任务入队
    for task in tasks:
        task_queue.put(task)
    
    def worker(processor):
        while True:
            # 阻塞取任务,队列空时线程休眠,无空转开销
            current_task = task_queue.get()
            try:
                processor(current_task)
            finally:
                # 无论任务执行是否报错,都标记任务完成,避免队列永久阻塞
                task_queue.task_done()
    
    # 每个处理器(对应一个节点)启动专属工作线程
    for p in processor_list:
        Thread(target=worker, args=(p,), daemon=True).start()
    
    # 阻塞等待所有任务处理完成
    task_queue.join()


if __name__ == "__main__":
    nodes = [Mork("Thad"), Mork("Chad"), Mork("Brad")]
    tasks = ["speak", "scream", "whisper", "speak", "speak", "whisper", "scream", "scream", "scream"]
    thread_pool = []

    for node in nodes:
        thread_pool.append(Processor(node))

    complete_all_tasks_using_threads(thread_pool, tasks)

关键特性

  • 无第三方依赖,全部基于Python标准库实现
  • 无额外锁开销:通过「单节点绑定单工作线程」的设计,从根源避免同一节点同时执行多个任务的冲突,不需要给节点实例加互斥锁
  • 任务领取逻辑线程安全:Queue的入队、出队操作都是原子性的,多线程并发领任务不会出现重复领取、丢任务的问题
  • 资源自动回收:工作线程设置为守护线程,所有任务执行完成后会随主线程自动退出,不需要手动维护线程生命周期
  • 异常安全:任务执行逻辑放在try/finally块中,即使任务执行抛错也会正确标记任务完成,不会导致调度逻辑永久卡死

内容的提问来源于stack exchange,提问作者Erik Lewis

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 06:48:23