跳转至

在基于线程的多 Agent 系统中,如何用 threading.Event、Lock、Queue 实现 Supervisor Worker 架构,保证任务分发的线程安全和优雅停止?

🎯 问题重述

在一个多线程 Agent 系统中,如何用 threading.EventLockQueue 搭建 Supervisor-Worker 架构,既保证任务分发线程安全,又能实现优雅停止?

核心目标:

  • 任务提交与消费互不阻塞

  • 多 Worker 竞争任务不出错

  • 停止信号下发后,Worker 能快速、安全地退出,不丢任务、不死锁


🧱 架构角色(用图标一目了然)
👑 Supervisor 线程
   ├── 创建任务 → 📮 Queue(线程安全任务缓冲)
   ├── 维护 🔒 Lock 保护的共享状态(如活跃任务数、停止标识)
   └── 收到停止信号 → 🚦 Event.set() 通知所有 Worker

👷 Worker 线程×N
   ├── 循环从 📮 Queue 获取任务(阻塞或超时)
   ├── 处理任务前/后更新共享状态(用 🔒 Lock)
   └── 检测到 🚦 Event 置位 → 处理完当前任务后安全退出

🧩 三个同步原语的分工

原语 角色 解决的核心问题
Queue 任务分发管道 自带线程安全,多生产者-多消费者无竞态
Lock 共享状态保护 保证“活跃任务计数”“是否正在停止”等变量的读写原子性
Event 停止信号广播 一个线程设置,所有等待/检查的线程立即感知,实现“一发多收”的优雅退出

💻 落地实现(面试可展示的关键代码)

import threading
import queue
import time
import random

class SupervisorWorkerSystem:
    def __init__(self, num_workers=3):
        self.task_queue = queue.Queue()          # 📮 任务队列
        self.lock = threading.Lock()             # 🔒 保护共享计数
        self.stop_event = threading.Event()      # 🚦 停止信号
        self.active_count = 0
        self.workers = []
        self.num_workers = num_workers

    def start(self):
        # 启动 Worker 线程
        for i in range(self.num_workers):
            t = threading.Thread(target=self._worker, args=(i,), daemon=False)
            self.workers.append(t)
            t.start()

    def submit_task(self, task_data):
        """Supervisor 提交任务(线程安全,Queue本身保证)"""
        with self.lock:
            self.active_count += 1
        self.task_queue.put(task_data)

    def _worker(self, worker_id):
        """Worker 主循环:取任务→处理→检查停止信号"""
        while not self.stop_event.is_set():
            try:
                # 阻塞等待任务,但设置超时以便周期性检查停止信号
                task = self.task_queue.get(timeout=1)
            except queue.Empty:
                continue   # 超时后回到循环首部检查 stop_event

            # 处理任务
            self._process_task(worker_id, task)

            # 任务完成,扣减计数
            with self.lock:
                self.active_count -= 1

    def _process_task(self, worker_id, task):
        """模拟任务处理"""
        print(f"👷 Worker-{worker_id} 处理任务: {task}")
        time.sleep(random.uniform(0.1, 0.3))   # 模拟耗时

    def graceful_shutdown(self):
        """Supervisor 触发优雅停止"""
        print("\n🛑 Supervisor 发送停止信号...")
        self.stop_event.set()                  # 通知所有 Worker

        # 等待 Worker 线程全部结束
        for t in self.workers:
            t.join()

        # 最后检查是否还有未完成任务(实际场景可记录日志)
        with self.lock:
            remaining = self.active_count
        print(f"✅ 所有 Worker 已退出,剩余未处理任务: {remaining}")

使用示例:

system = SupervisorWorkerSystem(num_workers=2)
system.start()

# Supervisor 不断投递任务
for i in range(5):
    system.submit_task(f"任务-{i}")
    time.sleep(0.2)

time.sleep(1)            # 等待部分任务执行
system.graceful_shutdown()

🔐 线程安全深度分析

  1. Queue 的天然安全性 task_queue.put()task_queue.get() 内部已加锁,多个 Supervisor 线程(如果有)和多个 Worker 同时操作队列,不会出现数据竞争或脏读。

  2. Lock 保护计数器的必要性 active_count 的修改(+=1-=1)是非原子操作,多线程同时执行会导致计数错误。通过 with self.lock 保证“读-改-写”互斥,确保任何时刻 Supervisor 看到的是真实任务存量。

  3. Event 的“无锁广播” stop_event 内部用条件变量实现,设置一次后所有线程的 is_set() 都立即可见,无需额外锁。它解决了“如何让阻塞在队列上的 Worker 同时醒来”的问题——结合 get(timeout=1),Worker 每隔1秒检查一次 Event,最差情况下1秒内也能感知到停止信号。


🛑 优雅停止机制的两个关键抉择

方案 A:毒丸(Sentinel) 向队列放入 N 个特殊值(如 None),Worker 收到后直接退出。 优点: 立即响应,不依赖超时。 缺点: 若队列中还有正常任务,毒丸会“插队”,导致任务丢失,需要额外逻辑处理。

方案 B:超时轮询 Event(推荐) Worker 用 get(timeout=T) 取出任务,每 T 秒检查一次 stop_event优点: 可以安全处理完队列中剩余任务再退出(只需在 stop_event 置位后继续取空队列的任务),且不会漏掉任何任务。 缺点: 最大延迟为 T 秒,通常设 0.5~1 秒对系统影响可忽略。

混合方式(最稳健): 设置 Event 后,Supervisor 仍然等待所有 Worker 自然消耗完队列,然后 join。Worker 在 stop_event 置位且队列为空时退出。这既保证了任务不丢,又避免了无限等待。


💡 进阶考量(让面试官点头的细节)

  • Worker 异常隔离:在 _process_task 外套 try...except,避免一个任务出错导致 Worker 线程崩溃。

  • 任务结果回传:可以再加一个 result_queue,Worker 处理完将结果 put 进去,Supervisor 异步收集。

  • 动态扩缩容:结合线程池(concurrent.futures)或根据 active_count 与队列长度动态增减 Worker。

  • 超时控制:可以为单个任务设置最大执行时间,避免某个任务拖死 Worker。


✅ 一句话总结

“Queue 解耦任务分发与处理,Lock 保证共享状态一致性,Event 实现零延迟广播停止信号——三者的组合让 Supervisor-Worker 架构天然线程安全,且具备生产级的优雅停止能力。”