线程通信与队列
约 1469 字大约 5 分钟
2026-05-10
多个线程之间需要传递数据时,不建议直接共享列表或字典。更推荐使用标准库 queue.Queue。
- 生产者消费者模型
- put 和 get
- taskdone 和 join
- 多个消费者
- 1写最简单的生产者消费者:q.put + q.get,比手动用 list + Lock 安全得多。
- 2用 q.put(item, timeout=1) 和 q.get(timeout=1) 处理“队列满/空”超时,避免线程死等。
- 3每个 q.get() 之后 finally 里写 q.task_done(),否则 q.join() 会永远卡住。
- 4N 个消费者要 put N 个哨兵 None;只 put 一个,只有一个消费者会退出。
- 5Queue(maxsize=100) 限积压量,防止生产者远快于消费者导致内存爆炸。
多个线程之间需要传递数据时,不建议直接共享列表或字典。更推荐使用标准库 queue.Queue。
Queue 是线程安全的队列,内部已经处理好锁,适合实现生产者消费者模型。
生产者消费者模型
生产者负责产生任务,消费者负责处理任务。
import queue
import threading
import time
task_queue = queue.Queue()
def producer():
for i in range(5):
task_queue.put(f'task-{i}')
print(f'生产任务 task-{i}')
def consumer():
while True:
task = task_queue.get()
print(f'处理任务:{task}')
time.sleep(1)
task_queue.task_done()
threading.Thread(target=consumer, daemon=True).start()
producer()
task_queue.join()
print('所有任务处理完成')put 和 get
put() 放入数据,get() 取出数据。
q = queue.Queue()
q.put('hello')
item = q.get()默认情况下:
- 队列满时,
put()会阻塞。 - 队列空时,
get()会阻塞。
可以设置超时:
try:
item = q.get(timeout=1)
except queue.Empty:
print('队列为空')task_done 和 join
当消费者处理完一个任务后,应该调用 task_done()。
主线程可以调用 join() 等待所有任务处理完成。
q.put('task')
item = q.get()
# 处理任务
q.task_done()
q.join()注意:每次 get() 后,都应该对应一次 task_done()。否则 join() 会一直等待。
多个消费者
import queue
import threading
import time
task_queue = queue.Queue()
def worker(worker_id):
while True:
task = task_queue.get()
print(f'worker-{worker_id} 处理 {task}')
time.sleep(1)
task_queue.task_done()
for i in range(3):
threading.Thread(target=worker, args=(i,), daemon=True).start()
for i in range(10):
task_queue.put(i)
task_queue.join()多个消费者可以并发处理队列中的任务。
哨兵值
如果不想使用 daemon 线程,可以用哨兵值通知消费者退出。
import queue
import threading
q = queue.Queue()
SENTINEL = object()
def worker():
while True:
item = q.get()
try:
if item is SENTINEL:
break
print(f'处理 {item}')
finally:
q.task_done()
threads = [threading.Thread(target=worker) for _ in range(3)]
for thread in threads:
thread.start()
for i in range(10):
q.put(i)
for _ in threads:
q.put(SENTINEL)
q.join()
for thread in threads:
thread.join()哨兵值是一个特殊对象,用来表示“没有更多任务了”。
Queue 的常见类型
| 类型 | 说明 |
|---|---|
queue.Queue | 普通先进先出队列 |
queue.LifoQueue | 后进先出队列,类似栈 |
queue.PriorityQueue | 优先级队列 |
queue.SimpleQueue | 简化队列,没有任务跟踪功能 |
大多数场景使用 Queue 即可。
注意事项
- 多线程传递任务时,优先使用
Queue。 - 每次
get()后要记得task_done()。 - 使用
join()等待队列任务全部完成。 - 需要停止消费者时,可以使用哨兵值。
- 不要依赖
qsize()做严格判断,它在多线程环境下只是近似值。
Queue 解决什么问题
多个线程之间最麻烦的是共享变量。queue.Queue 提供了一种更安全的方式:一个线程把任务放进去,另一个线程取出来处理。
你可以把它想成取餐口:
生产者:把订单放到取餐口
消费者:从取餐口拿订单处理生产者不需要知道消费者什么时候处理,消费者也不需要知道任务从哪里来。双方只通过队列交接。
一个更完整的生产者消费者例子
import queue
import threading
import time
q = queue.Queue()
def producer():
for i in range(5):
print('生产任务', i)
q.put(i)
time.sleep(0.2)
q.put(None) # 哨兵值,告诉消费者可以退出
def consumer():
while True:
item = q.get()
try:
if item is None:
print('消费者退出')
return
print('处理任务', item)
time.sleep(0.5)
finally:
q.task_done()
t1 = threading.Thread(target=producer)
t2 = threading.Thread(target=consumer)
t1.start()
t2.start()
q.join()
print('所有队列任务都处理完了')这里的 None 是哨兵值,意思是“没有更多任务了”。
多个消费者时哨兵值也要多个
如果有 3 个消费者,通常要放 3 个哨兵值。否则只有一个消费者收到退出信号,其他消费者还会继续等。
for _ in range(3):
q.put(None)这是初学者很容易漏掉的地方。
task_done 和 join 的配合
q.put(item):放入一个任务。q.get():取出一个任务。q.task_done():告诉队列“刚才那个任务处理完了”。q.join():等待所有放进去的任务都被task_done()确认。
如果调用了 get() 却忘记 task_done(),join() 可能一直等下去。
所以更稳的写法是:
item = q.get()
try:
handle(item)
finally:
q.task_done()Queue 不只是线程安全容器
Queue 的价值不只是“安全”,还在于让程序结构更清楚:
- 生产者只负责产生任务;
- 消费者只负责处理任务;
- 队列负责缓冲和交接;
- 可以通过
maxsize控制积压。
例如:
q = queue.Queue(maxsize=100)当队列满了,生产者会等待,这可以防止任务无限堆积导致内存暴涨。
总结
队列是线程之间传递数据的推荐方式。它比共享列表更安全,也更容易表达生产者消费者模型。复杂多线程程序中,尽量让线程通过队列通信,而不是共享变量。
- Queue 适合在线程之间传递任务,比共享列表和全局变量更安全清晰。
- 使用 `q.join()` 时,每次 `get()` 后都要配对 `task_done()`,最好放在 finally 里。
- 多个消费者需要多个哨兵值;`maxsize` 可以防止任务无限积压。
版权所有
版权归属:Shuo Liu
