如何基于自定义线程与对象实现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
相关产品推荐
相关产品推荐

