线程池:别一直雇佣厨师,叫外卖吧
开篇除恐
当你第一次看到程序里动态创建大量线程(比如为每个网络请求都 threading.Thread(target=task).start())时,可能会觉得“这不是很自然吗?有活儿就雇个人,忙完就解雇”。其实,频繁创建和销毁线程就像每来一单外卖就去招聘一个临时送餐员——招聘、培训、发工资、离职手续都要花时间和资源,而真正送餐的时间可能只占很小一部分。线程池正是为了解决这个问题:提前准备好一组固定数目的“空闲厨师”(线程),有任务来时直接分配给空闲的厨师去干,任务完成后厨师返回池中待命,这样就能省去大量线程创建/销毁的开销,同时还能控制并发度,防止资源耗尽。
白话话
- 线程池(Thread Pool):一个管理一组工作线程的容器。线程在池启动时就被创建并进入等待状态;任务(通常是一个可调用对象)被提交到池中的工作队列;空闲线程从队列中取任务并执行;执行完后再次返回队列等待下一个任务。
- 工作队列(Work Queue):存放待处理任务的队列(常见为先进先出)。生产者(提交任务的线程)把任务放入队列;消费者(工作线程)从队列取任务。
- 同步机制:为了安全地把任务放入队列和从队列取任务,需要互斥锁保护队列的访问,并常配合条件变量实现“队列为空时工作线程等待”“队列满时生产者等待”(若有上限)。
- 生命周期管理:线程池需要优雅地关闭:先停止接受新任务,然后通知所有工作线程在完成当前任务后退出;最后等待所有工作线程结束(join)。
- 参数:
pool_size:同时工作的线程数量,通常设置为 CPU 核心数或根据任务类型(I/O 密集型还是 CPU 密集型)适当调大。max_queue_size(可选):队列的最大长度,用于防止无限制增长导致内存耗尽;当队列满时,新任务的提交者可以被阻塞或拒绝。
- 优点:
- 减少线程创建/销毁开销;
- 控制并发数,避免系统因过多线程而导致上下文切换开销升高;
- 简化编程模式:提交任务即可,无需手动管理线程生命周期。
- 注意事项:
- 若任务本身会阻塞很长时间(如网络 I/O),线程池大小可能需要调大,以免所有线程都被阻塞导致队列积压;
- 任务函数应尽量不要捕获易失效的引用(如栈局部变量的引用),除非确保任务在引用生命周期结束前完成;
- 异常处理:任务中抛出的异常应被捕获并记录,否则可能导致工作线程意外退出。
直觉先行
想象你经营一家快餐店: - 每来一单外卖,你不必去街边招聘一个临时送餐员、培训他怎么骑车、发工资、再办离职手续;而是店里一直有固定的几位送餐员在待命。 - 有单子来时,你就叫一个空闲的送餐员:“去把这单送到XX路。”送餐员完成后返回店里,继续待命下一单。 - 这样,无论是早上的高峰还是下午的低谷,你都能保证有足够的送餐员来处理订单,而不会因为招聘和离职频繁而浪费时间和金钱。
在编程里,线程池就是那些一直在待命的送餐员,工作队列就是外卖订单列表,提交任务就是下单。只要线程池的大小设得合理,你就能够在不频繁创建/销毁线程的情况下高效处理大量短任务。
例子贴身
例子 1(☼ 热身):使用 concurrent.futures.ThreadPoolExecutor 简单模拟线程池
任务:演示如何使用标准库的线程池来并发执行多个短任务,并获取结果。
怎么想到的:我们准备 10 个任务,每个任务只是打印自己的 ID 并休眠 50 ms。我们用 ThreadPoolExecutor 提交它们,存放返回的 Future,最后遍历 future 调用 result() 等待完成。
示例代码(Python)
import concurrent.futures
import time
def task(task_id):
print(f"Task {task_id} started")
time.sleep(0.05) # 50 ms
print(f"Task {task_id} finished")
return task_id * 2
def main():
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
futures = [executor.submit(task, i) for i in range(10)]
results = [f.result() for f in futures]
print("All tasks done")
print("Results:", results)
if __name__ == "__main__":
main()
运行结果(任务交叉进行)
Task 0 started
Task 1 started
Task 2 started
...
Task 0 finished
Task 1 finished
...
All tasks done
Results: [0, 2, 4, 6, 8, 10, 12, 14, 16, 18]
要点:
- ThreadPoolExecutor 在创建时就预先创建好指定数量的工作线程,之后只负责调度任务。
- 使用 submit 提交任务并返回 Future 对象,调用者可以通过 result() 获取返回值或捕获异常。
- 上下文管理器 with 确保池在使用完后自动关闭(调用 shutdown(wait=True)),等待所有线程完成。
- 真正的线程池会复用线程,避免每次任务都产生新线程的开销。
例子 2(☼☼ 正经):手写简单线程池(固定大小,无界队列)
任务:实现一个最小的线程池类 ThreadPool,支持:
- 构造时指定线程数量;
- 提交任务(任何可调用对象,返回 future 以获取结果或仅执行);
- 安全关闭(shutdown() 方法不再接受新任务,待所有任务完成后让工作线程退出)。
怎么想到的:我们内部使用:
- list 存放工作线程;
- queue.Queue 作为任务队列;
- threading.Lock 和 threading.Condition 实现工作线程在队列为空时等待;
- atomic 风格的布尔标志(用 threading.Event 或简单的变量配合锁)标记池是否已停止接受新任务;
- 工作线程的循环:等待有任务或止步信号;取出任务并执行;如果是止步且队列空则退出。
示例代码(Python)
import threading
import queue
import time
class ThreadPool:
def __init__(self, num_workers):
self._tasks = queue.Queue()
self._workers = []
self._stop_event = threading.Event()
self._lock = threading.Lock() # 用于保护 _stop_event 读取(虽然 Event 本身是线程安全的,这里为了演示锁的使用)
for _ in range(num_workers):
worker = threading.Thread(target=self._worker_loop)
worker.start()
self._workers.append(worker)
def _worker_loop(self):
while True:
task = None
with self._lock:
if self._stop_event.is_set() and self._tasks.empty():
break
try:
# 非阻塞获取任务,若无任务则短暂等待
task = self._tasks.get(timeout=0.1)
except queue.Empty:
continue
if task is not None:
try:
task()
finally:
self._tasks.task_done()
def submit(self, func, *args, **kwargs):
"""提交一个任务,返回一个 Future-like 对象(这里仅演示,不返回真实 future)"""
future = threading.Event()
result = [None]
exc = [None]
def wrapper():
try:
result[0] = func(*args, **kwargs)
except Exception as e:
exc[0] = e
finally:
future.set()
self._tasks.put(wrapper)
return _SimpleFuture(future, result, exc)
def shutdown(self, wait=True):
"""停止接受新任务,并可选择等待所有线程结束"""
self._stop_event.set()
if wait:
for w in self._workers:
w.join()
class _SimpleFuture:
def __init__(self, event, result_container, exc_container):
self._event = event
self._result = result_container
self._exc = exc_container
def result(self):
self._event.wait()
if self._exc[0]:
raise self._exc[0]
return self._result[0]
def main():
pool = ThreadPool(4) # 4 条工作线程
fut1 = pool.submit(lambda: time.sleep(0.1) or 42)
fut2 = pool.submit(lambda x: x * x, 7)
print("Result1:", fut1.result()) # 42
print("Result2:", fut2.result()) # 49
pool.shutdown(wait=True)
if __name__ == "__main__":
main()
运行结果
Result1: 42
Result2: 49
要点:
- 线程池在构造时就创建好所有工作线程,之后只负责调度任务。
- 使用自定义的 _SimpleFuture 模拟了 future 接口,实际项目中可直接使用 concurrent.futures.Future 或更完整的实现。
- shutdown 方法负责优雅关闭:先设置停止标志,唤醒所有工作线程(通过队列超时或事件),让它们在看到停止且队列空时退出,然后 join 每个线程。
- 这个实现为无界队列;若想加上队列长度上限,只需在 submit 前检查队列大小并在达到上限时阻塞或抛出异常。
例子 3(☼☼☼ 硬骨头):使用 concurrent.futures 与 threading.Barrier 实现分阶段同步(概念演示)
任务:展示如何结合线程池和同步原语来完成更复杂的协作,比如在所有工作线程完成某个阶段后才进入下一个阶段。
怎么想到的:threading.Barrier 可以让一组线程在达到某个点之前相互等待。我们可以在线程池中提交一批任务,每个任务在完成第一阶段后在 barrier 上等待,所有任务到达后同时继续第二阶段。这里仅作概念展示,不给出完整代码,而是说明思路。
要点:
- concurrent.futures.ThreadPoolExecutor 与 threading.Barrier, threading.Event, threading.Condition 等同步原语配合使用,可以构建丰富的线程协作模式。
- 实际开发中,推荐使用成熟的线程池实现(如 concurrent.futures.ThreadPoolExecutor、multiprocessing.pool.ThreadPool、第三方库如 pebble、asyncio 等),除非有特定定制需求。
收尾
这一章要带走的东西
- 线程池的核心思想是:预先创建固定数目的工作线程,复用它们来执行提交的任务,从而避免频繁线程创建/销毁的开销。
- 线程池主要由 工作队列 + 互斥锁 + 条件变量 组成;工作线程在队列为空时等待,生产者在提交任务时唤醒一个工作线程。
- 使用 future/promise(如 concurrent.futures.Future)可以让任务提交既简单又安全,支持获取返回值和捕获异常。
- 优雅的关闭流程是:先停止接受新任务(设置标志),唤醒所有工作线程,让它们在完成当前任务且队列为空时退出,最后 join 所有线程。
- 实际开发中,推荐使用成熟的线程池实现(如 concurrent.futures.ThreadPoolExecutor、multiprocessing.pool.ThreadPool、第三方库如 ctpl、ThreadPool 等),除非有特定定制需求。
- 掌握线程池后,你就能够在高并发场景(如网络服务器、任务调度、批处理)中高效地调度计算资源,而不必担心线程爆炸带来的系统开销。
- 下一章我们将给出新手上路的检查清单:从创建第一个线程到使用线程池,涵盖最常见的坑和最佳实践,帮助你在实际项目中安全、高效地使用多线程。
就这样。 下一章我们来看看“新手上路的检查清单”,把前面所学的知识点提炼成实用的步骤和注意事项,让你在写多线程代码时少走弯路。