ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python多线程实战:从GIL到线程池的完整指南

Python多线程实战:从GIL到线程池的完整指南 很多 Python 初学者学到多线程时都会抱着这样一个期待开了多个线程程序执行速度应该能翻倍吧然后写一段死循环测试却发现结果出乎意料——不仅没变快有时候反而更慢了。于是网上开始流传另一种声音Python 多线程就是假的完全没有用别学了。这两种说法其实都只对了一半。Python 的threading模块不是万能的加速器但也不是一个该被抛弃的鸡肋。它真正擅长解决的是 I/O 密集型任务比如网络请求、文件读写、数据库查询而在 CPU 密集型计算场景下它会受到 GIL全局解释器锁的限制。如果你正在写爬虫、做批量接口调用、处理大量文件或者想在 GUI 程序里避免界面卡死threading都是非常实用的工具。这篇文章不从概念堆砌开始而是直接告诉你多线程适合解决什么问题、不适合解决什么问题然后从threading模块的核心 API 出发用代码演示三种创建线程的方式、线程锁、事件、队列、线程池以及实际工程中常见的坑和排查思路。读完你可以照着代码跑一遍并且能判断你自己的项目到底该不该用多线程。1. 为什么你写的 Python 多线程没有变快先回到开头的困惑为什么多线程没有让代码更快这里必须先解释一个 Python 特有的机制——GIL全称 Global Interpreter Lock也就是全局解释器锁。CPython 解释器在同一个进程内同一时刻只允许一个线程执行 Python 字节码。这意味着如果你有 8 个线程在跑 CPU 密集型计算它们并不会真正地同时运行在 8 个 CPU 核心上而是轮流抢同一把锁。# 一个直观的类比 单线程进程食堂只有一个打饭窗口 多线程进程食堂还是只有一个打饭窗口只是排队的人分成 8 列 多进程方案开了 8 个窗口每列人同时打饭所以如果是纯计算任务比如循环求和、图片像素处理、大数据量的数学运算多线程反而会因为线程创建、上下文切换、锁竞争带来额外开销表现不如单线程。这种场景下真正合适的是multiprocessing多进程或者直接用 NumPy 这类能把计算下沉到 C 层的库。但 I/O 密集型任务就完全不同。当线程执行requests.get()、file.read()、time.sleep()这类操作时线程会进入等待状态此时 GIL 会被释放其他线程就能继续执行。也就是说多线程在等待网络响应或磁盘读写的时间里可以让其他线程去干活从而实现“并发等待”整体耗时会大幅缩短。概括成一张判断表任务类型特点推荐方案原因I/O 密集型网络请求、文件读写、数据库查询threading 或 asyncio等待期间自动让出 GIL并发效果明显CPU 密集型数值计算、循环处理、压缩解压multiprocessing绕开 GIL真正利用多核混合型既有计算又有等待threading 计算下沉到 C 库提高吞吐计算部分交给更合适的工具这篇文章后面说的所有示例默认都是面向 I/O 密集型场景。2. threading 模块的核心概念与适用场景2.1 线程与主线程一个 Python 程序启动后默认会有一个执行流这个执行流叫主线程。使用threading模块创建的额外执行流称为子线程。每个线程都有自己独立的方法调用栈但共享同一个进程内的全局变量和内存空间。这是多线程和多进程最大的差异多线程之间通信成本低但需要处理共享数据的竞争问题。看一个最简单的最小示例# 文件路径demo_basic_thread.py import threading import time def worker(): print(f子线程开始: {threading.current_thread().name}) time.sleep(1) print(f子线程结束: {threading.current_thread().name}) if __name__ __main__: print(f主线程开始: {threading.current_thread().name}) t threading.Thread(targetworker, nameworker-thread) t.start() t.join() print(f主线程结束: {threading.current_thread().name})运行结果主线程开始: MainThread 子线程开始: worker-thread 子线程结束: worker-thread 主线程结束: MainThread这里有两个容易被忽略的细节第一target指定线程要执行的函数name给线程命名方便日志排查。第二t.join()会让主线程等待子线程执行完毕再继续。如果去掉join()主线程不会等子线程程序可能直接执行完退出子线程还没来得及完整输出。初学者可以先把线程理解成“程序里另开的几条小流水线”。主线程负责统筹子线程负责干活join()就是“等这条流水线收工”。2.2 守护线程线程有一个布尔属性daemon中文叫守护线程。daemonTrue的线程不会阻止主线程退出。主线程结束时守护线程会被强制终止。默认情况下threading.Thread创建的非守护线程主线程会等待它执行完才退出。# 文件路径demo_daemon.py import threading import time def background_task(): while True: print(后台任务运行中...) time.sleep(0.5) if __name__ __main__: t threading.Thread(targetbackground_task, daemonTrue) t.start() time.sleep(2) print(主线程即将退出守护线程会被强制结束)注意守护线程通常用于日志轮转、心跳检测、自动保存这类“后台服务型”任务。业务中重要的任务不要设置成守护线程否则主线程退出后任务可能执行到一半被强杀。2.3 线程生命周期一个线程从创建到结束大致经历以下状态状态说明触发方式创建Thread 对象被实例化threading.Thread(targetfn)就绪/运行线程开始执行t.start()阻塞/等待线程等待 I/O 或锁遇到time.sleep()、queue.get()死亡函数执行完毕或发生未捕获异常自动结束或异常退出threading模块没有直接提供获取线程运行状态的方法但可以通过t.is_alive()判断线程是否还在运行。3. 环境准备与实验说明本文的所有代码基于 Python 3.10 版本编写但涉及的 API 在 Python 3.6 之后的版本都能正常使用。如果你的电脑还没装 Python或者不确定当前版本可以先执行下面的命令检查python --version建议在项目目录下创建一个独立的虚拟环境避免依赖包污染系统环境# Windows python -m venv venv venv\Scripts\activate # Linux / macOS python3 -m venv venv source venv/bin/activate为了让后续的示例效果更明显建议安装requests库用于网络请求场景的演示pip install requests如果你使用的是 Anaconda它自带requests直接跳过这一步即可。本文的核心代码不依赖第三方库的只有前三节从第 5 节网络请求示例开始才需要requests。4. 三种创建线程的常见方式很多教程一上来就讲多线程的完整架构其实没必要。你先掌握三种创建线程的方式后续再谈复杂的线程协作。4.1 方式一直接通过 Thread 类传入函数这是最常用、也最推荐新手使用的方式。# 文件路径demo_create_thread.py import threading import time def fetch_url(url): print(f正在请求: {url}) time.sleep(1) print(f请求完成: {url}) if __name__ __main__: urls [ https://example.com/api/1, https://example.com/api/2, https://example.com/api/3, ] threads [] for url in urls: t threading.Thread(targetfetch_url, args(url,)) threads.append(t) t.start() for t in threads: t.join() print(所有请求已完成)关键点在于args(url,)。注意这里必须是一个元组如果写成args(url)就等于传入一个字符串Python 会报TypeError。当只有一个参数时一定要记得加结尾逗号。4.2 方式二继承 threading.Thread 重写 run 方法这种方式适合业务逻辑比较复杂、需要在线程中保存状态的场景。# 文件路径demo_subclass_thread.py import threading import time class DownloadThread(threading.Thread): def __init__(self, task_id): super().__init__() self.task_id task_id self.result None def run(self): print(f任务 {self.task_id} 开始) time.sleep(1) self.result f任务 {self.task_id} 的返回值 print(f任务 {self.task_id} 结束) if __name__ __main__: threads [] for i in range(1, 4): t DownloadThread(i) threads.append(t) t.start() for t in threads: t.join() for t in threads: print(t.result)这段代码里有一个非常实用的设计把执行结果保存在self.result上。join()之后主线程可以从线程对象中取出子线程的运行结果。这比使用全局变量来收集结果要清晰得多。继承方式的好处是代码组织性好每个线程相当于一个具有内部状态的任务对象缺点是写法略冗余。简单任务用方式一就够了。4.3 方式三ThreadPoolExecutor 线程池每次任务都手动创建线程在线程数量大时效率不高。更工程化的做法是使用线程池预先创建一批线程任务来了就分配一个空闲线程执行。concurrent.futures.ThreadPoolExecutor是 Python 内置的线程池实现# 文件路径demo_thread_pool.py from concurrent.futures import ThreadPoolExecutor, as_completed import time def handle_task(task_id): time.sleep(1) return f任务 {task_id} 处理完成 if __name__ __main__: task_ids list(range(1, 11)) with ThreadPoolExecutor(max_workers3) as executor: future_map {executor.submit(handle_task, task_id): task_id for task_id in task_ids} for future in as_completed(future_map): task_id future_map[future] try: print(future.result()) except Exception as e: print(f任务 {task_id} 执行失败: {e})运行结果顺序可能不同任务 1 处理完成 任务 3 处理完成 任务 2 处理完成 任务 5 处理完成 ... 任务 10 处理完成这里解释几个常用方法executor.submit(fn, *args)提交一个任务到线程池返回一个Future对象。future.result()获取任务返回值如果任务抛异常这里会重新抛出异常。as_completed(futures)是一个迭代器哪个任务先完成就先返回哪个 future适合处理执行时间不定的任务。with语句退出时会调用executor.shutdown()等待所有已提交任务执行完毕。线程池真正解决了“线程数量管理”的问题。如果不限制线程数你就创建 1 万个线程去请求接口操作系统资源会被快速耗尽。合理设置max_workers可以保护系统资源。5. 线程同步多线程不是各跑各的线程之间共享全局变量这会带来一个经典问题多个线程同时修改一个数据最终结果不符合预期。5.1 不加锁的竞态条件先看一个经典的计数器示例# 文件路径demo_race_condition.py import threading counter 0 def increment(): global counter for _ in range(1000000): counter 1 if __name__ __main__: threads [] for _ in range(10): t threading.Thread(targetincrement) threads.append(t) t.start() for t in threads: t.join() print(fcounter 的最终值: {counter})如果程序按顺序正确执行10 个线程各自加 100 万次结果应该是 10000000。但实际运行多次结果往往是几百万、几十万每次都不一样。原因在于counter 1并不是一个原子操作。它内部可以拆成三步读取 counter 当前值。计算 counter 1。把新值写回 counter。假设两个线程同时读取到 counter 100线程 A 算出了 101线程 B 也算出了 101然后各自写回counter 最终只变成 101而不是 102。这就丢失了一次更新。5.2 加锁解决竞态条件threading.Lock可以保证同一时刻只有一个线程能进入临界区# 文件路径demo_lock.py import threading counter 0 lock threading.Lock() def increment(): global counter for _ in range(1000000): with lock: counter 1 if __name__ __main__: threads [] for _ in range(10): t threading.Thread(targetincrement) threads.append(t) t.start() for t in threads: t.join() print(fcounter 的最终值: {counter})这次每次运行结果都应该是 10000000。with lock:是 Python 上下文管理器语法等价于lock.acquire() try: counter 1 finally: lock.release()第二种写法也不能算错但with语法更安全。如果临界区中间的代码抛出异常finally还能保证锁被释放直接调用acquire()后忘记release()就会造成死锁。5.3 可重入锁 RLockLock 还有一个明显的限制同一个线程不能连续多次acquire()。如果代码里存在嵌套加锁或者一个方法内部调用另一个也加锁的方法就可能死锁。这时需要使用threading.RLock可重入锁。它允许同一个线程多次获取锁内部用一个计数器记录获取次数只在全部release()后才真正释放锁。# 文件路径demo_rlock.py import threading lock threading.RLock() counter 0 def inner(): with lock: global counter counter 1 def outer(): with lock: print(外层加锁) inner() print(内层调完counter , counter) if __name__ __main__: outer()如果这里用的是普通Lock在outer()中获取锁后再调用inner()尝试获取同一把锁程序会直接卡死因为同一线程无法获取自己已经持有的 Lock。RLock就能避免这个问题。5.4 实际场景多线程写入同一个文件写日志是多线程中很常见的场景。如果多个线程同时往一个文件里写会出现日志内容交错、行被拆断的问题。给写入操作加一把写锁是最直接的解决方案。# 文件路径demo_file_write.py import threading import time write_lock threading.Lock() def write_log(filename, thread_name, content): time.sleep(0.01) with write_lock: with open(filename, a, encodingutf-8) as f: f.write(f[{thread_name}] {content}\n) if __name__ __main__: threads [] for i in range(10): t threading.Thread(targetwrite_log, args(app.log, fT{i}, fline {i})) threads.append(t) t.start() for t in threads: t.join() print(日志写入完成)这里的time.sleep(0.01)故意模拟日志内容的生成耗时。加上写锁后同一时刻只有一个线程在操作文件写入内容就是完整的行。6. 线程之间的通信与协作锁解决了“数据竞争”但很多场景下线程之间还需要互相通知和传递数据。6.1 使用 Queue 安全传递任务最推荐的线程间通信方式是使用queue.Queue。它内部已经实现了线程安全put()和get()方法都自带了锁机制不需要你再去加一把锁。# 文件路径demo_queue.py import queue import threading import time import random task_queue queue.Queue() result_queue queue.Queue() def producer(): for i in range(10): task f任务-{i} task_queue.put(task) print(f生产: {task}) time.sleep(random.uniform(0.1, 0.3)) def worker(worker_id): while True: try: task task_queue.get(timeout3) except queue.Empty: break print(f线程 {worker_id} 处理: {task}) time.sleep(random.uniform(0.2, 0.5)) result_queue.put(f{task} 已完成) task_queue.task_done() if __name__ __main__: producer_thread threading.Thread(targetproducer) producer_thread.start() workers [] for i in range(3): t threading.Thread(targetworker, args(i,)) workers.append(t) t.start() producer_thread.join() for t in workers: t.join() task_queue.join() print(所有任务处理完毕)这段代码模拟了一个典型的生产者消费者模型生产者线程往task_queue里放任务。多个工作线程从队列里取任务并处理。处理结果放入result_queue之后你可以用另一个消费者线程来收集结果。关键方法说明方法作用put(item)向队列添加元素如果队列满则阻塞get()从队列取出元素如果队列空则阻塞get_nowait()非阻塞获取元素空队列时抛出queue.Emptytask_done()告诉队列当前任务已处理完join()阻塞直到队列里所有任务都被标记为task_doneget(timeout3)是为了避免工作线程在队列清空后一直阻塞。当队列空了且 3 秒内没有新任务线程会捕获queue.Empty后退出。如果不设超时get()会永远阻塞下去程序无法正常结束。6.2 使用 Event 做线程间通知threading.Event适合从一个线程向另一个线程发送“信号”。它内部维护一个布尔标志set()会把它变为 Truewait()会阻塞直到标志变为 True。# 文件路径demo_event.py import threading import time start_event threading.Event() def worker(worker_id): print(f线程 {worker_id} 已就绪等待启动信号...) start_event.wait() print(f线程 {worker_id} 收到信号开始工作) if __name__ __main__: threads [] for i in range(3): t threading.Thread(targetworker, args(i,)) threads.append(t) t.start() time.sleep(1) print(主线程发送启动信号) start_event.set() for t in threads: t.join()运行结果线程 0 已就绪等待启动信号... 线程 1 已就绪等待启动信号... 线程 2 已就绪等待启动信号... 主线程发送启动信号 线程 0 收到信号开始工作 线程 2 收到信号开始工作 线程 1 收到信号开始工作Event 常用于以下场景多个工作线程初始化后先等待主线程统一发出“开始”信号。工作线程需要周期性等待某个外部配置更新的信号。需要优雅关闭线程时通过 Event 通知循环线程退出。7. 完整案例用多线程并发请求接口把前面的知识点综合起来做一个贴近真实工作的案例批量请求一组 URL并使用线程池管理并发度最后返回每个请求的状态和耗时。# 文件路径demo_concurrent_requests.py import threading import time import requests from concurrent.futures import ThreadPoolExecutor, as_completed # 仅用于演示实际请求需要保证目标域名允许访问 URLS [ https://www.example.com, https://www.python.org, https://docs.python.org, ] lock threading.Lock() results [] def fetch_one(url): start time.time() try: resp requests.get(url, timeout5) status resp.status_code length len(resp.text) except Exception as exc: status str(exc) length 0 cost round(time.time() - start, 2) # 多线程同时向列表追加数据可能造成竞争这里加一个锁保护 with lock: results.append((url, status, length, cost)) return url, status, cost if __name__ __main__: start_all time.time() with ThreadPoolExecutor(max_workers3) as executor: futures [executor.submit(fetch_one, url) for url in URLS] for future in as_completed(futures): url, status, cost future.result() print(f完成: {url}, 状态: {status}, 耗时: {cost}s) total_cost round(time.time() - start_all, 2) print(f总耗时: {total_cost}s) print(\n结果汇总:) for url, status, length, cost in results: print(f{url} - 状态码: {status}, 内容长度: {length}, 耗时: {cost}s)注意requests.get()是阻塞式网络请求。当线程发出网络请求后它会在等待响应时释放 GIL所以多线程能够明显缩短整体请求时间。如果所有请求是串行的N 个请求的总耗时约为 N 个请求耗时之和使用线程池后总耗时约等于最慢的那个请求耗时前提是线程池大小足够。实际工程中线程池的max_workers建议结合目标服务器的承受能力来设置并不是越大越好。大量并发请求可能触发对方限流或造成 IP 被临时封禁。合规爬虫应当控制请求频率并遵守目标网站的robots.txt和访问协议。8. 常见问题与排查思路多线程程序的调试难度比单线程高很多。遇到问题不要盲改代码先按下面的表格逐项排查。问题现象可能原因排查方式解决方案程序运行结束后线程还没执行完主线程未调用join()检查代码中是否对每一个Thread调用了join()在main最后统一join()所有线程多个线程同时修改同一个列表/字典数据丢失共享数据写入存在竞态条件在写入位置打印线程名和当前数据长度使用threading.Lock保护写入或改用queue.Queue程序卡死无法退出死锁或某个线程阻塞在get()上使用py-spy dump查看线程栈或打印每个线程的启动状态避免嵌套锁使用RLock给queue.get()设置超时CPU 使用率始终只有 100% 左右纯 CPU 计算任务受 GIL 限制观察任务类型是计算密集还是 I/O 密集改用multiprocessing或把计算逻辑改用 NumPy线程函数抛异常但主线程看不到子线程异常默认不会传播到主线程在run()方法外层捕获异常并打印 traceback自定义线程基类统一捕获并记录异常线程池任务执行顺序和提交顺序不一致并发任务本来就是无序完成打印各任务完成时间需要严格有序时改用单线程或按索引收集结果daemonTrue的线程任务没执行完程序就退出守护线程不阻止主线程退出确认线程任务是否具备事务性非守护任务不要设置daemonTruemax_workers设置很大但速度没有提升I/O 服务和本地资源存在瓶颈检查目标服务响应时间、本机文件句柄数适当降低线程数或调整系统文件句柄限制再补充两个排查工具和技巧第一使用threading.enumerate()打印当前存活线程列表import threading for thread in threading.enumerate(): print(f线程名: {thread.name}, 是否存活: {thread.is_alive()})第二使用 Python 自带的日志模块记录线程名。logging默认支持%(threadName)s格式能在日志中直接看到是哪条线程打印的信息import logging logging.basicConfig( levellogging.INFO, format%(asctime)s [%(threadName)s] %(message)s ) logging.info(这条日志会显示线程名)排查死锁时这两个信息通常能快速定位问题。9. 最佳实践与工程建议多线程代码写起来容易写对很难。这里总结几条我在实践中认为最重要的工程建议。优先使用线程池而不是手动创建大量线程。ThreadPoolExecutor帮你管理线程生命周期减少线程创建和销毁的开销同时能限制最大并发数量。手动创建线程适合线程数量很少、场景简单的程序。共享数据结构统一通过加锁或 Queue 访问。 不要把“往字典里加一条数据”当成一个绝对安全的小操作。Python 对于单个赋值操作在 CPython 下往往带有 GIL 保护但复合操作如list.append前先判断长度仍然可能出错。最稳妥的方式是所有跨线程的共享状态变更都加锁或者使用线程安全的queue.Queue。锁的粒度要尽可能小。 锁住的代码范围越小线程等待时间越短并发效率越高。比如有 1000 万次循环不要在循环外面一次性加锁然后循环 1000 万次那样等于把多线程又变成串行了。正确做法是每次操作前加锁操作完立刻释放。不过频繁加锁也会带来性能消耗需要在实践中找到平衡。使用超时和取消机制。 队列取任务的get(timeout...)、等待锁的acquire(timeout...)、等待线程结束的join(timeout...)这些超时参数在异常场景下非常重要。生产环境中一个线程可能因为外部服务无响应而卡死超时机制能保证程序不会永久阻塞。异常必须在线程内部捕获。 线程函数里抛出的异常不会自动传递给主线程。如果你想知道哪个线程发生了什么错误必须在子线程内部捕获并记录。更好的做法是在线程函数的入口写统一的try/except并调用日志模块。不要为了用多线程而用多线程。 任务总数很少比如只有两个请求单线程串行可能只需要 0.5 秒用线程池反而要花时间创建线程。使用多线程前先量化任务的等待时间和数量。如果任务本身很快并发带来的收益可以忽略不计。注意线程本身的启动开销。 线程创建不是零成本的。如果任务是百万级别的短任务每条线程又很轻量更适合使用线程池复用线程。线程数量超过几千以后线程调度本身会成为瓶颈。10. 总结与后续学习方向回到最开始的问题。Python 的threading模块不是让你把 CPU 计算加速到多核并行的银弹它是 Python 在 I/O 密集型场景下提升吞吐量的实用工具。理解 GIL 的边界是理解 Python 多线程的前提。在 I/O 等待场景下threading可以让多个请求并发等待大幅缩短总耗时在 CPU 密集计算场景下应该选择multiprocessing或把核心计算交给 C 扩展库。本文通过完整代码演示了三种创建线程的方式、Lock 和 RLock 的用法、Queue 如何做线程间通信、Event 如何传递信号以及线程池在生产场景中的标准写法。这些知识点放在一个真实的批量请求任务中串起来之后你已经具备在 Python 爬虫、后端服务、自动化脚本里接入多线程的基本能力。下一步你可以从两个方向继续深入一个是向异步编程方向探索Python 的asyncio在单线程内使用事件循环处理大量连接在高并发 I/O 场景下是另一种解决方案。多线程和多进程适合并发执行多个独立任务而异步适合大量任务之间快速切换。另一个是向多进程方向探索学习multiprocessing模块、进程池、进程间通信理解 Python 中进程与线程的选择边界。建议你不只是在本地运行示例代码而是找一个真实的小项目做实验。找一个需要并发请求的接口列表用单线程、多线程、线程池分别测一遍总耗时和资源占用你就对 Python 并发模型有了更直观的体感。欢迎收藏这篇文章后续遇到线程相关问题可以回来对照排查表快速定位。
返回列表