线程池:别一直雇佣厨师,叫外卖吧

开篇除恐

当你第一次看到程序里动态创建大量线程(比如为每个网络请求都 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.Lockthreading.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.futuresthreading.Barrier 实现分阶段同步(概念演示)

任务:展示如何结合线程池和同步原语来完成更复杂的协作,比如在所有工作线程完成某个阶段后才进入下一个阶段。

怎么想到的threading.Barrier 可以让一组线程在达到某个点之前相互等待。我们可以在线程池中提交一批任务,每个任务在完成第一阶段后在 barrier 上等待,所有任务到达后同时继续第二阶段。这里仅作概念展示,不给出完整代码,而是说明思路。

要点
- concurrent.futures.ThreadPoolExecutorthreading.Barrier, threading.Event, threading.Condition 等同步原语配合使用,可以构建丰富的线程协作模式。
- 实际开发中,推荐使用成熟的线程池实现(如 concurrent.futures.ThreadPoolExecutormultiprocessing.pool.ThreadPool、第三方库如 pebbleasyncio 等),除非有特定定制需求。

收尾

这一章要带走的东西
- 线程池的核心思想是:预先创建固定数目的工作线程,复用它们来执行提交的任务,从而避免频繁线程创建/销毁的开销。
- 线程池主要由 工作队列 + 互斥锁 + 条件变量 组成;工作线程在队列为空时等待,生产者在提交任务时唤醒一个工作线程。
- 使用 future/promise(如 concurrent.futures.Future)可以让任务提交既简单又安全,支持获取返回值和捕获异常。
- 优雅的关闭流程是:先停止接受新任务(设置标志),唤醒所有工作线程,让它们在完成当前任务且队列为空时退出,最后 join 所有线程。
- 实际开发中,推荐使用成熟的线程池实现(如 concurrent.futures.ThreadPoolExecutormultiprocessing.pool.ThreadPool、第三方库如 ctplThreadPool 等),除非有特定定制需求。
- 掌握线程池后,你就能够在高并发场景(如网络服务器、任务调度、批处理)中高效地调度计算资源,而不必担心线程爆炸带来的系统开销。
- 下一章我们将给出新手上路的检查清单:从创建第一个线程到使用线程池,涵盖最常见的坑和最佳实践,帮助你在实际项目中安全、高效地使用多线程。

就这样。 下一章我们来看看“新手上路的检查清单”,把前面所学的知识点提炼成实用的步骤和注意事项,让你在写多线程代码时少走弯路。