ARTICLE DETAIL

资讯详情

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

Python多线程实战:Thread类、GIL与线程池核心机制

Python多线程实战:Thread类、GIL与线程池核心机制 写 Python 写久了你会发现很多程序的性能瓶颈其实不在 CPU 算不快而是卡在“等待”上——等网络响应、等磁盘读写、等数据库返回。这时候Thread 类也就是threading标准库里的Thread就成了顺手就能用的并发工具。它允许你同时跑多个任务让等待时间里 CPU 去做别的事整体吞吐量能明显提升。这篇文章我会用实际跑过的代码和经验把Thread类的基本用法讲透从最基础的创建和启动到join、守护线程、锁、队列再到线程池的取舍和线上排查技巧。适合刚接触多线程、或者用得不深想系统捋一遍的开发者。看完你就能直接上手也知道哪些坑打死不能踩。1. Thread 到底在解决什么问题1.1 先从程序执行的本质说起一个普通 Python 程序默认是从上到下一行一行执行的——这叫单线程模型。遇到input()、time.sleep()、requests.get()这类会阻塞的操作时程序就停在那里干等什么都不干。做过爬虫的人肯定深有体会同一个 URL 列表挨个请求一个响应等几百毫秒几百个 URL 就要好几分钟实际上 CPU 利用率极低。Thread 类解决的问题简单说就是把“等待”的时间利用起来。你开一条线程去发网络请求主线程继续做其他事请求回来之后再让另一段代码处理结果。这样一来程序从“串行排队”变成“多个队列同时推进”整体的墙钟时间实际跑完所花的真实时间会大幅缩短。要注意这里的用词多线程是并发不一定是并行。并发是“看起来同时在做”多个任务交替执行并行才是“真正同时执行”需要多核 CPU 配合。Python 标准库里的Thread在这件事上有自己的特殊性下面单独说。1.2 Thread 类能解决什么问题Thread 的典型应用场景我归纳成三类IO 密集型的提速网络请求、文件读写、数据库查询、消息队列消费。这类任务大部分时间在等外部系统多线程能把等的时间重叠起来。后台任务的托管比如程序启动后需要持续监听一个端口、轮询一个配置文件、定期上报心跳用守护线程跑后台逻辑主程序继续做自己的事。并发的业务场景比如同时下载多个文件、批量处理多个用户请求、模拟并发压测一个接口。它不适合做的一件事是纯计算密集型的任务比如大规模矩阵运算、复杂的加密哈希。这种任务几乎不等待一直在吃 CPU。由于 Python 的全局解释器锁GILGlobal Interpreter Lock限制同一时刻只有一个线程能执行 Python 字节码开线程不但不会加速反而可能因为上下文切换变慢。如果你真有这种需求应该考虑多进程multiprocessing或者用 C 扩展、异步 IO。1.3 线程不等于并行GIL 的影响GIL 是 Python 标准实现CPython里的关键机制。它保证同一时刻只有一个原生线程能跑 Python 字节码。很多初学者听到这个会很困惑那线程到底有什么用关键在于阻塞操作。当线程在等网络数据、读磁盘、调sleep的时候它其实是在操作系统层面挂起的这时候 GIL 会被释放其他线程可以继续跑 Python 代码。所以 GIL 掐掉的是“同时使用多个 CPU 核心的能力”并没有掐掉“并发处理 IO 的能力”。这也是我用“Thread 类来提升程序吞吐”的前提——你的任务构成里等待占比要高。如果任务本身既有计算又有等待也可以考虑组合方案主计算逻辑用进程池进程内的 IO 部分用线程。这个后面章节结合线程池再展开。2. Thread 类核心用法从创建到销毁2.1 最简单的创建与启动方式先上一个最小示例这是所有 Thread 用法的地基import threading import time def worker(name: str): print(f线程 {name} 开始工作) time.sleep(2) print(f线程 {name} 工作结束) t threading.Thread(targetworker, args(下载任务,)) t.start() print(主线程继续执行其他工作) t.join() print(所有任务完成)这段代码做了这么几件事threading.Thread(targetworker, args(下载任务,))创建了一个线程对象target是线程要执行的函数args是传给该函数的参数元组。t.start()让线程被操作系统调度从这一刻起worker开始执行。主线程没有傻等而是先打印了“主线程继续执行其他工作”。t.join()让主线程等t跑完再继续。这里有一个新手常犯的错把start()写成了run()。函数对象本身有个run()方法如果直接调用t.run()它不是在线程里执行而是在当前线程同步执行效果跟没开线程一样。如果你看到程序完全按顺序跑、看不到并发效果第一件事就查是不是调了run()。2.2 target 与 args 的细节target参数除了传普通函数还能传类实例的可调用方法比如传t threading.Thread(targetobj.fetch_page, args(url,))这个对象的fetch_page会被调用。lambda 表达式适用于简短逻辑比如t threading.Thread(targetlambda: print(hello))。不推荐传复杂逻辑调试不方便。任何有__call__方法的对象类实例如果实现了__call__也能直接作为 target。args必须是一个元组。如果你只传一个参数很容易忘记加逗号(url)其实是一个字符串(url,)才是元组。传字典参数用kwargst threading.Thread(targetworker, kwargs{name: 下载任务})再强调一点线程的执行顺序是没法预料的。start()的调用顺序不代表实际执行顺序这取决于操作系统的调度策略。如果你想控制先后就得用线程锁或者Event这个后面说。2.3 join 到底在等什么join()这个方法的语义是调用它的线程通常是主线程会阻塞直到被join的那个线程终止。光看这句话可能觉得很简单但它有几个容易忽略的细节join 可以传超时参数t.join(timeout5)表示最多等 5 秒超时就不等了继续往下走。这在主程序需要快速退出或不需要等全部完成时很有用。要注意超时返回时被 join 的线程可能还在后台跑。join 一个已经结束的线程不会报错会立刻返回。join 不能对一个还没 start 的线程调用会直接抛RuntimeError: cannot join thread before it is started。如果你开了一批线程想等它们全部结束用列表管理它们threads [] for i in range(10): t threading.Thread(targetworker, args(i,)) t.start() threads.append(t) for t in threads: t.join()上面这种“先全部 start 再全部 join”的写法是有讲究的。如果写成交替 start-join效果会退化成串行——开一个线程等它结束再开下一个。所以批量开线程时一定先把所有线程都start()起来最后再统一join()。2.4 daemon 线程后台任务的生命周期管理daemon 这个词直译是“守护”你可以把守护线程理解成“跟随主线程生死的后台任务”。设置方式很简单t threading.Thread(targetheartbeat_check, daemonTrue) t.start()默认情况下创建的线程是非守护的daemonFalse。这会导致一个问题如果主线程执行完所有代码程序不会立即退出而是等所有非守护线程结束才退出。这确实是多线程程序经常遇到的“程序卡着退不出去”的罪魁祸首之一。你启动了几个非守护线程它们因为等某个外部资源一直阻塞着主线程早就跑完了程序却一直挂着。如果把线程设为 daemonTrue主线程结束运行时这些守护线程会被强制终止程序就能正常退出。典型的用途包括后台轮询、心跳上报、监控日志、临时监听端口。但反过来你也得小心使用守护线程不保证代码完整执行。比如守护线程正在写一个文件写到一半主线程退出这个文件就是坏的。涉及数据一致性、必须完整落盘的任务千万不能用守护线程而应该用队列配合 join 的方式优雅关闭。场景daemon 设置原因主程序退出后无所谓的后台统计任务daemonTrue避免拖住进程退出必须完整执行完的日志落盘任务daemonFalse保证数据不丢网络服务端的工作线程daemonFalse每个连接处理必须结束定时清理临时文件的任务daemonTrue进程退出时没必要继续3. 线程安全与同步绕不开的坎3.1 为什么多线程会出乱子多线程最有价值的一点是共享数据——多个线程读写同一个变量就能协作。但协作的前提是大家都按规矩来否则就会产生竞态条件。举个实际例子。一个抢购程序里tickets是总票数多个下单线程执行“先检查库存再扣减库存”的逻辑import threading tickets 20 def buy(): global tickets if tickets 0: # 模拟网络延迟和判断处理 import time time.sleep(0.01) tickets - 1 threads [threading.Thread(targetbuy) for _ in range(50)] for t in threads: t.start() for t in threads: t.join() print(f余票: {tickets}) # 结果往往不是 0可能是 30、40 甚至更多结果往往是负数或者没有扣减到底。原因在于if tickets 0和tickets - 1不是原子操作线程 A 检查到票数大于 0还没来得及扣减线程 B 也检查到同样结果两个线程就都执行了扣减本来库存只有 20实际被扣了更多次。这种场景下必须用同步机制。Thread类本身没有提供任何数据保护能力标准做法是用Lock、RLock、Event或者Queue。3.2 Lock最简单的互斥保护Lock就像一把门锁线程进入“临界区”前要拿到锁用完再还回来。拿到锁的线程在锁释放前其他线程都会阻塞在acquire()上。改造上面的例子import threading tickets 20 lock threading.Lock() def buy(): global tickets with lock: if tickets 0: time.sleep(0.01) tickets - 1用with lock而不是手动lock.acquire()/lock.release()好处是即使临界区里抛了异常锁也会被正确释放不会出现死锁。这里要注意一个 Lock 的特性它是不能重复持有的。同一线程如果对同一个 Lock 调用了两次acquire()会直接把自己阻塞住程序就卡死了。例如lock threading.Lock() with lock: with lock: # 第二次获取时永远等不到因为锁被自己持有 pass这种场景需要用RLock可重入锁。RLock允许同一个线程多次获取锁内部的计数器会记录持有次数只有释放次数匹配才算真正释放。递归函数、嵌套代码结构里经常出现这种需求。我还想补充一个 Lock 的原理细节阻塞线程把自己挂起之前会发生操作系统层面的线程上下文切换。这个开销并不小。如果你在临界区里做非常耗时的操作所有其他线程只能干等效率反而很差。所以锁的粒度要尽量小——只保护真正需要原子的那几行代码不要在整个函数外面套一个大锁。3.3 用 Event 做协作控制有些场景不需要争抢资源而是一个线程等待另一个线程发出“可以开始”的信号。比如主线程让所有工作线程同时开始压测、或者通知某个线程停止运行。start_event threading.Event() stop_event threading.Event() def worker(): print(等待开始指令...) start_event.wait() # 阻塞到事件被 set while not stop_event.is_set(): print(干活中) time.sleep(1) t threading.Thread(targetworker, daemonTrue) t.start() time.sleep(2) start_event.set() # 放行 time.sleep(3) stop_event.set() # 通知停止Event的机制非常朴素内部维护一个标志位set()把它置为真clear()把它置为假wait()在标志为假时阻塞为真时立刻返回。你不需要管底层怎么调度这个工具在设计上就是为了线程间发信号。这里有一个特别常用的工程模式用Event优雅地停止一个后台线程。线程里循环判断“stop 事件是否被设置”主程序准备退出时set()一下线程在完成当前迭代后自然退出。这比直接强制杀线程安全得多。3.4 用 Queue 在线程间安全传数据线程之间传递数据最稳妥的方式不是共享一个变量并手动加锁而是用queue.Queue。Queue内部已经实现好了线程安全的锁和条件变量你只管put和get就行。典型的生产者-消费者模型长这样import queue import random q queue.Queue(maxsize10) def producer(): for i in range(100): item f任务-{i} q.put(item) # 队列满时会阻塞 print(f产生 {item}) time.sleep(random.random()) def consumer(): while True: item q.get() # 队列空时会阻塞 if item is None: # 结束信号 break print(f消费 {item}) q.task_done() threads [ threading.Thread(targetproducer), threading.Thread(targetconsumer, daemonTrue), ] for t in threads: t.start() q.join() # 等待所有任务被处理完Queue有几个好用的特性maxsize控制队列长度满了之后put会阻塞天然形成背压不会让内存无限制上涨。get()默认阻塞配合timeout参数可以在空队列超时后做其他逻辑。用发送None或别的哨兵值来做“优雅结束”信号消费者拿到哨兵就退出。这是做关闭控制最常用的手段。queue.Queue的底层实现是deque和threading.Condition你不需要自己加锁。实战中一个线程池配一个任务队列是比裸开无数个Thread更稳的架构。这块内容放在第 4 章讲。4. 线程池与 Thread 类的取舍4.1 频繁创建线程的代价有人会想既然 Thread 这么好用我每来一个任务就创建一个新线程行不行短时间来看没问题任务量一大你就会发现问题。创建线程本身是有代价的每次创建线程操作系统都需要分配内核栈、用户态栈建立线程控制块分配 TLS线程局部存储。如果任务执行时间只有几毫秒而线程创建耗时占了大头程序整体吞吐反而下降。另外一个程序如果同时开着几千个线程光上下文切换就能把 CPU 全部吃满任务还得排队。我在压测一个爬虫系统时就踩过这个坑用 Thread 每个 URL 开一个线程跑到 200 个并发时程序开始疯狂 GC垃圾回收爬到一半直接 OOM。后来把方案改成线程池加有限队列问题立刻消失。4.2 ThreadPoolExecutor 怎么接替 ThreadPython 标准库提供了concurrent.futures.ThreadPoolExecutor本质上就是线程池的官方实现。from concurrent.futures import ThreadPoolExecutor, as_completed def fetch_url(url): return requests.get(url).status_code urls [fhttps://example.com/page/{i} for i in range(20)] with ThreadPoolExecutor(max_workers8) as executor: future_to_url {executor.submit(fetch_url, url): url for url in urls} for future in as_completed(future_to_url): url future_to_url[future] try: status future.result() print(f{url} - {status}) except Exception as e: print(f{url} 请求失败: {e})核心思路是max_workers控制同时运行的线程数量。设置多少合适我的经验是CPU 核心数 * (5-10)针对 IO 密集任务但也要结合下游系统的承受能力比如目标接口的 QPS 上限。submit返回一个Future对象可以理解为“未来才完成的结果”。你可以之后调用future.result()取真实返回或者捕获异常。用with语句会自动等待所有任务完成再关闭线程池。这段代码如果按照上一章的思路去套相当于线程池帮我管理了创建、调度、等待、回收这些脏活。ThreadPoolExecutor还有一个兄弟叫ProcessPoolExecutor底层的线程换成了进程。遇到纯计算密集型任务直接换成ProcessPoolExecutor就能绕开 GIL 的限制。但这意味着参数和返回值需要被序列化用 pickle传输成本会高一些。4.3 什么场景必须用 Thread 类线程池确实好用但它不给线程命名、不好设置 daemon、也不好针对每个线程单独做生命周期管理。有些场景你还是得亲手用Thread类需要长时间运行的常驻线程比如一个消费者线程持续监听任务队列从启动一直活到进程退出。这种线程用线程池去托管反而别扭。需要精细控制单个线程的启动顺序比如某个任务必须在另外两个任务成功启动之后再开始Thread 对象可以被个别地 start 和 join。需要设置线程名称、自定义 run 行为继承 Thread 类重写run()方法在线程里初始化自己的上下文这些线程池很难做得这么细。你完全可以组合使用线程池负责日常突发的短任务手动创建的 Thread 负责长期驻留的后台任务。5. 常见问题与排查技巧实录5.1 明明开了多线程时间却没什么变化这是大家最容易遇到的一段爬虫代码开 100 个线程结果耗时只比单线程快 5%。排查思路看一眼任务是不是 CPU 密集。如果是循环计算、哈希、图像处理多线程受 GIL 限制不该有明显提升。换成多进程试一下。看是不是有全局锁在排队。把关键函数里的with lock打印日志看看竞争是不是很激烈。看下游阻塞时间是否不足以抵消线程开销。如果每个请求只要 1 毫秒线程切换成本可能比收益还大。一个实用的判断方法把任务函数里的 IO 部分用time.sleep(0.1)模拟跑一次单线程再跑一次 10 线程对比耗时。如果提速接近 10 倍说明场景适合多线程如果只有一两倍就要重新设计了。5.2 程序退不出去卡在 join 上现象是主线程代码执行完了程序却一直挂在那不动。最常见的原因就是存在非守护线程一直没结束。排查流程打印threading.enumerate()看当前还有哪些线程活着。这个函数会列出所有 Thread 对象包括主线程和后台线程。逐个检查这些线程在等什么可能是从Queue.get()取数据时队列一直空可能是Lock.acquire()拿不到锁也可能是在等某个网络超时。确定哪些线程是应该存在的、哪些是泄漏的。泄漏线程通常在请求处理逻辑里不小心 new 了 Thread但没有加回收机制。该设 daemon 的设daemonTrue该用“哨兵消息join”模式处理的就发结束信号。如果你只是想“不卡死”把不会正常结束的线程设成守护线程主线程退出时会直接终止它们。但注意这可能导致未写完的数据丢失工程上还是优先做优雅关闭。5.3 线程里变量被改得乱七八糟你以为只在函数里用了局部变量结果多个线程互相干扰。我单独强调一次Python 的局部变量是线程私有的全局变量是线程共享的类实例属性也是线程共享的。如果你在多个线程里操作同一个对象的同一个属性必须加锁或者用threading.local()。threading.local()的设计是给每个线程独立的数据空间import threading context threading.local() def set_user(name): context.user name def show_user(): print(f当前线程 {threading.current_thread().name} 的 user 是 {context.user})浏览器 Web 框架里的“请求上下文”就是这么实现的。每个请求线程只能看到自己存的属性天然隔离。如果你发现一个线程改了数据导致其他线程异常先检查是不是忘了用local()。5.4 线程抛了异常却没人知道线程函数里抛出的异常不会自动传到主线程如果没人捕获只会打印一段堆栈然后线程悄悄死掉。排查这种问题很费劲因为主线程完全没有感知。我习惯在写线程入口函数时统一包一层 try-except把异常记录到日志或者放入结果队列def safe_worker(func, *args, **kwargs): try: return func(*args, **kwargs), None except Exception as e: return None, e如果是ThreadPoolExecutor的submit方案异常会存在Future里调future.result()时会重新抛出这个体验好很多。另外给线程设置唯一的namethreading.Thread(nameorder-queue-consumer)在排查时也特别有用日志里能直接看到是哪个线程出了问题。线上排查线程还常用到threading.active_count()配合监控看线程数量是否异常增长用faulthandler.dump_traceback_later(10)把当前所有线程的堆栈快照打出来特别适合抓“线程卡死在哪一行”。6. 多线程编程的几点工程经验最后分享几个我长时间写多线程程序的实际心得。第一把线程当作一种受控资源而不是免费午餐。我见过最夸张的代码每个 HTTP 请求进来都threading.Thread(...).start()服务直接被打垮。线程数量的上限取决于你的 IO 模型、内存大小和下游系统能力不是越大越好。用线程池或者信号量做限量是长期可靠的方案。第二能用队列就尽量别共享变量。多线程协作最稳的模式是“数据通过队列流动”而不是两个线程直接改同一个列表。队列天然有锁和通知机制出错少、也容易调试。你只需要把注意力放在队列里的数据格式上。第三关闭程序要有预案。多线程程序最难的部分其实不是启动而是关闭。老是感叹“程序退不出去了”。我现在的习惯是每个长期线程至少监听一个Event或队列哨兵主线程退出时先发关闭信号再join(timeout5)超时了再选择强制结束。这样既保证数据落盘又不会无限期挂起。第四多线程代码一定要加日志和名字。线程数量一多没有名字的线程出错时只能看到Thread-1、Thread-13这样的编号完全没法对应业务。我给自己定的规矩所有线程创建时必须传name日志里用threading.current_thread().name记录来源。排查效率立竿见影。Thread 类是 Python 并发编程的基础砖块你只需要把第 2 章的创建方式、第 3 章的同步机制、第 4 章的资源管理策略吃透大部分业务场景都能稳稳拿捏。真遇到需要更高吞吐或复杂协作的时候再往 asyncio 和多进程方向拓展也不迟。
返回列表