在基于线程的多 Agent 系统中,如何用 threading.Event、Lock、Queue 实现 Supervisor Worker 架构,保证任务分发的线程安全和优雅停止?
🎯 问题重述¶
在一个多线程 Agent 系统中,如何用
threading.Event、Lock、Queue搭建 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()
🔐 线程安全深度分析¶
-
Queue 的天然安全性
task_queue.put()和task_queue.get()内部已加锁,多个 Supervisor 线程(如果有)和多个 Worker 同时操作队列,不会出现数据竞争或脏读。 -
Lock 保护计数器的必要性
active_count的修改(+=1、-=1)是非原子操作,多线程同时执行会导致计数错误。通过with self.lock保证“读-改-写”互斥,确保任何时刻 Supervisor 看到的是真实任务存量。 -
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 架构天然线程安全,且具备生产级的优雅停止能力。”