Skip to content

线程与线程池实战:把等待时间叠起来

对应视频:EP04《线程与线程池实战》
视频入口:B站 EP04 定时稿(计划公开:2026-08-11 09:54)。 对应合集:Python 并发实战:从基础模型到语音工程。 本文为视频的工程展开版,补充代码模板、排查链路、状态字段和检查清单。 最后核验:2026/08/09 08:26:33

视频解释线程池适合等待型任务。文章给出 ThreadPoolExecutor 的实战模板:控制并发数、收集结果、处理异常,避免无限提交。

线程池任务生命周期

这张图把线程池拆成提交、等待、收敛、异常和复盘五步;线程池的难点不在启动线程,而在受控地结束每个任务。

最小模板

python
from concurrent.futures import ThreadPoolExecutor, as_completed


def handle_one(item: str) -> tuple[str, int]:
    # 替换成真实 I/O:请求接口、读取文件、调用外部服务等。
    return item, len(item)


def run_batch(items: list[str], max_workers: int = 8) -> list[tuple[str, int]]:
    results: list[tuple[str, int]] = []
    errors: list[tuple[str, Exception]] = []

    with ThreadPoolExecutor(max_workers=max_workers) as pool:
        future_to_item = {pool.submit(handle_one, item): item for item in items}

        for future in as_completed(future_to_item):
            item = future_to_item[future]
            try:
                results.append(future.result())
            except Exception as exc:
                errors.append((item, exc))

    if errors:
        for item, exc in errors[:5]:
            print(f"failed: {item}: {exc}")
        raise RuntimeError(f"{len(errors)} tasks failed")

    return results

什么时候适合线程池

场景是否适合说明
批量请求 HTTP 接口适合等网络,线程可推进其它任务
读取大量小文件适合但要限并发受磁盘能力影响
调用外部命令并等待适合但要限并发同时受 CPU、磁盘和进程数影响
纯 Python 计算通常不适合默认 CPython 下会受 GIL 影响
写同一个结果文件不直接适合需要单写入者或锁

工程补充:怎么把这篇用到项目里

读图方式

  • 提交任务前先确定输入规模和 max_workers,不要把几十万任务一次性 submit 进内存。
  • as_completed 逐个收敛结果,这样失败可以早发现、早归档。
  • 线程池结束后一定复盘吞吐、错误率、资源水位,再决定是否调大或调小并发。

排查路径

  • 用 1、4、8、16 个 worker 跑同一批输入,观察吞吐是否线性提升、错误是否集中出现。
  • 如果 CPU 低但吞吐提升明显,说明等待可叠加;如果 CPU 或连接数先满,要停下来加限制。
  • 失败任务要记录输入、异常类型和 traceback,避免只在控制台看到一堆 Future 报错。

问题、证据和解决动作

问题现场先查什么解决动作
一次性提交全量文件RSS、pending future 数分批 submit 或用有界队列
某个请求永久卡住单任务 elapsed、无 timeout给外部调用加 timeout
异常被吞掉结果数量少但程序正常退出调用 future.result() 并记录 future.exception()
线程越多越慢连接池、磁盘 util、错误率降低 max_workers,按资源拆池

落地边界

  • 线程池适合等待型任务,不适合拿来加速纯 Python CPU 密集循环。
  • 所有外部调用都要有超时,否则一个 worker 可以无限占住池容量。

常见失败模式

错误写法现象修正
不读取 future.result()异常被藏在 Future 里as_completed 收集结果和异常
max_workers 随便开很大内存、连接数、磁盘压力上升从小并发压测,观察瓶颈
所有任务一次性提交大批任务常驻内存分批提交或使用队列
多线程直接写同一文件结果错乱或丢行单写入者模式

实战检查清单

  • 是否有 max_workers 上限。
  • 是否读取了每个 Future 的结果。
  • 是否记录失败任务。
  • 是否避免多个线程直接写同一份结果。
  • 是否用顺序版本做过基准对照。

读完之后能完成什么

  • 能写出可收集异常的线程池模板。
  • 能判断线程池适不适合当前任务。
  • 能避免把线程池当作无限加速按钮。

配套示例代码

源码文件:examples/ep04/threadpool_template.py

python
from concurrent.futures import ThreadPoolExecutor, as_completed


def handle_one(item: str) -> tuple[str, int]:
    return item, len(item)


def run_batch(items: list[str], max_workers: int = 8) -> list[tuple[str, int]]:
    results: list[tuple[str, int]] = []
    errors: list[tuple[str, Exception]] = []

    with ThreadPoolExecutor(max_workers=max_workers) as pool:
        future_to_item = {pool.submit(handle_one, item): item for item in items}

        for future in as_completed(future_to_item):
            item = future_to_item[future]
            try:
                results.append(future.result())
            except Exception as exc:
                errors.append((item, exc))

    if errors:
        for item, exc in errors[:5]:
            print(f"failed: {item}: {exc}")
        raise RuntimeError(f"{len(errors)} tasks failed")

    return results


if __name__ == "__main__":
    print(run_batch(["a.wav", "b.wav", "c.wav"], max_workers=2))
别急,先让缓存热一下。