ARTICLE DETAIL

资讯详情

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

Python多线程threading模块实战:掌握GIL、锁与队列

Python多线程threading模块实战:掌握GIL、锁与队列 在 Python 开发里多线程是一个绕不开的话题。接触这个主题最常见的方式就是学习 threading 模块创建线程、等待线程结束、在多个线程之间传递任务、解决共享变量的竞争问题。很多人学完 Thread 创建和 start 之后就以为掌握了多线程实际一写业务脚本就发现要么结果不对要么线程没有退出要么程序直接卡住。这是按照 100 天精通 Python 系列第 37 天的内容整理的 threading 模块基础与实战。适合已经会 Python 基本语法、但还不太清楚线程之间如何协作的读者。学完以后你可以解释 GIL 对多线程的影响用 threading.Thread 创建线程用 join 控制主线程退出用 Lock 保护共享数据用 queue.Queue 分发任务并收集结果遇到线程没有执行、结果丢失、程序挂起时能按路径排查。1. 先厘清线程、进程与 GIL再决定要不要用多线程1.1 进程与线程隔离、共享与切换成本进程是操作系统分配资源的基本单位每个进程拥有独立的地址空间。进程之间不能直接访问对方内存必须用管道、消息队列、共享内存等机制通信。好处是一方崩溃不容易影响另一方坏处是创建和通信成本高。线程是在进程内运行的执行单元。同一进程下的多个线程共享代码段、数据段和堆内存但各自拥有独立的栈和寄存器上下文。线程的创建成本比进程低线程之间天然适合共享数据这是 Python 多线程经常被用来并发的原因。放到实际场景里如果一个脚本需要同时读取多个文件或者同时请求多个下游服务多线程是一个合理的实现方式。因为每个线程只需要做好自己的任务操作系统会负责线程调度。多线程的核心难点不是创建线程而是处理共享数据时的并发访问问题。多个线程同时读没有问题但只要有一个线程要写共享状态就要考虑原子性、可见性和加锁代价。1.2 GIL 的作用与反直觉结论GIL 全称是 Global Interpreter Lock也就是全局解释器锁。CPython 是 Python 最常见的官方实现这个实现里的解释器同一时刻只允许一个线程执行 Python 字节码。因为 GIL 的存在很多人会得出“Python 多线程没用”的结论这个说法并不完整。关键点是GIL 限制的是一段 Python 字节码的执行不是限制操作系统线程本身。当线程执行到 IO 操作时比如 sleep、文件读写、网络请求、数据库等待解释器会释放 GIL让其他线程去占用解释器。所以多线程对 IO 密集任务通常有效因为程序的大量时间都在等待外部资源。真正受 GIL 影响明显的是 CPU 密集任务。如果代码在持续做数值计算、字符串哈希、大循环两个线程争夺同一个 GIL无法同时使用多个 CPU 核心结果往往是速度不升反降。即使在多核机器上CPython 的默认解释器行为也是这样。这里不要和“Python 不可能并发”混淆。Python 可以在一个进程里创建大量线程系统层面线程也确实由操作系统调度只是在同一时刻解释器只会执行某一个线程的 Python 代码。对于很多脚本线程仍然能让整体耗时大幅下降因为占用的是等待时间段。1.3 适用场景速查IO 密集、CPU 密集与混合任务任务类型典型例子多线程效果更合适的方案IO 密集文件读写、sleep、HTTP 请求、数据库查询等待效果明显线程等待时释放 GILthreading、concurrent.futures.ThreadPoolExecutorCPU 密集大数乘法、加密散列、纯 Python 循环计算效果差受 GIL 限制难并行multiprocessing、ProcessPoolExecutorIO 密集 少量 CPU 计算下载后解析处理、读取文件后按行统计可以拆分处理仍建议固定 worker 线程queue threading或进程池配合线程池高并发、大量长连接上千个 socket 连接同时等待数据线程数量过多会带来切换成本asyncio 协程实际工程里任务往往不是单纯的 IO 或 CPU。可以先写一个最简版本跑通再用注释标记哪些步骤真正阻塞。如果阻塞来自外部资源多线程或异步是合理的如果阻塞来自 CPU 计算就要考虑多进程或改用带 native 线程释放 GIL 的库。2. 用 Thread 写第一个多线程脚本并学会控制线程生命周期2.1 最小示例创建线程并观察执行顺序threading.Thread 是创建线程最直接的类。target 参数传入一个可调用对象start 方法启动线程线程会在自己的执行上下文里运行 target。下面是一个最基础的多线程例子import threading import time def worker(): for i in range(3): time.sleep(0.5) print(fworker 运行中: {i}, flushTrue) t threading.Thread(targetworker) t.start() print(主线程继续执行, flushTrue)运行这段代码时输出顺序通常不会稳定。主线程继续执行可能会先打印也可能在 worker 执行到一半时打印。原因是 start 之后主线程和子线程并行运行谁先拿到 CPU 时间片由操作系统决定。这是多线程入门的第一课不要依赖线程执行顺序。如果代码逻辑依赖顺序需要在代码里加同步机制而不是期望线程按创建顺序执行。2.2 用 join 保证线程执行完成如果主线程需要等待子线程执行完再继续就调用 join。join 的含义是“当前线程等待目标线程结束”。import threading import time def worker(): for i in range(3): time.sleep(0.5) print(fworker 运行中: {i}, flushTrue) t threading.Thread(targetworker) t.start() t.join() print(worker 已经执行完主线程继续, flushTrue)加入 join 之后主线程会阻塞在这里直到 t 里的 worker 返回。只有确认线程处理完毕主线程才继续执行。join 还可以传 timeout 参数例如t.join(timeout2)。含义是主线程最多等两秒超时后即使子线程还没结束join 也会返回。这对防止程序无限等待很有用但要注意超时返回不代表子线程已经退出后续流程要处理“任务还没有完成”的场景。2.3 传参、命名与把结果带回来Thread 支持通过 args 和 kwargs 向目标函数传参。线程的名字在生产环境调试时非常重要建议为每个线程显式命名日志里就能区分是哪个线程在处理任务。import threading import time def read_resource(resource_name, timeout): print(f{threading.current_thread().name} 开始读取 {resource_name}, flushTrue) time.sleep(timeout) return f{resource_name} 的内容 t threading.Thread( targetread_resource, args(本地配置文件,), kwargs{timeout: 1}, nameconfig-reader, ) t.start() t.join()需要特别说明的是Thread 不能像普通函数一样直接接收返回值。如果要把一个函数的返回值带回主线程通常需要借助共享容器。最简单的方式是传入一个 list 或 dict让线程把结果放进去。import threading results [] def worker_with_result(resource_name): result f{resource_name} 处理完成 results.append(result) t threading.Thread( targetworker_with_result, args(订单数据,), ) t.start() t.join() print(results)这里直接把结果 append 到 list在简单示例里没有问题。一旦多个线程同时写同一个 list 或 dict就要考虑加锁或者直接改用 queue.Queue后文的第 3 章和第 5 章会展开讲。2.4 为什么会遇到“线程根本没有执行”这是初学者最常见的报错式困惑t threading.Thread(targetworker())这里把worker()的返回值传给了 target。Python 会先普通地调用一次 worker然后把返回值交给线程对象看起来就像“线程没有执行但函数执行过了”。正确写法是传函数对象而不是调用结果t threading.Thread(targetworker)如果 worker 需要参数就用 args 传t threading.Thread(targetworker, args(任务A,))排查时不要只看代码位置。可以在目标函数里打印线程名确认 start 是否真的把函数调度到新线程。线程启动只调用一次 start重复调用会抛出 RuntimeError。3. 当多个线程操作同一份数据Lock 是怎么起作用的3.1 一个看起来正常但结果不对的计数器多线程访问共享变量最经典的例子是计数器。import threading count 0 def add(): global count for _ in range(1000000): count 1 threads [threading.Thread(targetadd) for _ in range(10)] for t in threads: t.start() for t in threads: t.join() print(最终结果:, count)理论上10 个线程各自加 100 万次结果应该是 1000 万。实际运行后结果往往小于 1000 万并且每次运行结果可能不同。原因在于count 1并不是一步完成的操作。它在 Python 字节码层面至少包含读取 count、计算加一、写回 count 三个步骤。线程 A 可能先读取 count还没写回线程 B 也读取了同一个旧值。于是两个线程都基于同一个值计算导致一次更新被覆盖。这个例子未必每次运行都错误但循环次数越大、线程越多出现竞争的概率越高。可以把循环次数调到 200 万再试。重点不是某个运行结果而是操作本身不具备原子性。3.2 用 with lock 修好计数器Lock 是线程同步的一个基本工具。获得锁的线程可以进入临界区其他线程必须等它释放后才能进入。import threading count 0 lock threading.Lock() def add(): global count for _ in range(1000000): with lock: count 1 threads [threading.Thread(targetadd) for _ in range(10)] for t in threads: t.start() for t in threads: t.join() print(加锁后的结果:, count)with lock会在进入代码块时 acquire离开代码块时自动 release即使代码块里抛异常也会释放锁。这比手动写 acquire 和 release 更加安全。Lock 解决的是互斥问题即“同一时间最多一个线程操作共享资源”。代价是其他线程必须等待因此锁的粒度要尽可能小。不要用一个大锁包住一大段代码否则多线程的并发优势会被抵消变成近似串行执行。3.3 不要用普通 Lock 代替 RLock普通 Lock 是不可重入锁。同一个线程如果已经持有锁再次 acquire 同一个 Lock会出现死锁。import threading lock threading.Lock() def outer(): with lock: inner() def inner(): with lock: print(inner) t threading.Thread(targetouter) t.start() t.join()outer 获取锁之后调用 innerinner 又尝试获取同一把 Lock这个线程会永远等自己释放锁程序就会卡住。这种情况可以用可重入锁 RLock 替代。RLock 允许同一个线程多次获取锁内部会记录持有次数每次 acquire 对应一次 release只有全部释放后其他线程才能获取。import threading lock threading.RLock() def outer(): with lock: inner() def inner(): with lock: print(inner) t threading.Thread(targetouter) t.start() t.join()如果方法之间互相调用并且都依赖同一把锁优先考虑 RLock不要靠“拆成两把锁”硬绕。多把锁反而可能引入死锁风险。3.4 一个更省心的设计直接用 queue 传递任务和结果多线程共享可变数据时很多问题都可以通过“减少共享”来规避。如果业务是任务分发模型优先用 queue.Queue。queue.Queue 内部已经用锁和通知机制实现了线程安全。put 往队列放数据get 从队列取数据多个线程同时读写同一个 Queue 时不需要再额外加锁。import queue import threading task_queue queue.Queue() def worker(name): while True: item task_queue.get() if item is None: task_queue.task_done() break print(f{name} 处理 {item}, flushTrue) task_queue.task_done() for i in range(5): task_queue.put(f任务-{i}) t threading.Thread(targetworker, args(worker-1,)) t.start() task_queue.join() task_queue.put(None) t.join()使用 queue 的思路是线程不依赖全局列表而是通过队列接口传递不透明的数据包。这样既避免自己写锁也让任务边界更清晰。需要注意get 之后要记得调 task_done否则主线程的 join 永远等不到队列清空。这个细节是后面实战章节最容易埋坑的地方。4. 线程之间的通信与并发控制Event、Semaphore 和 daemon 线程4.1 用 Event 通知其他线程开始Event 适合一个线程通知其他线程“现在可以继续”的场景。Event 内部维护了一个 bool 标志初始为 False。wait 的线程会阻塞直到其他线程调用 set 把标志置为 True。import threading import time event threading.Event() workers [] def wait_and_run(name): print(f{name} 已就绪等待开始信号, flushTrue) event.wait() print(f{name} 收到信号开始工作, flushTrue) for i in range(3): t threading.Thread(targetwait_and_run, args(fworker-{i},)) t.start() workers.append(t) time.sleep(1) event.set() for t in workers: t.join()Event 适合做“闸门”比如先让所有线程初始化完成再统一放行。和 join 的区别在于Event 并不等待线程结束它只是在线程之间传递一个状态。如果初始化完成后不再需要通知可以调用 event.set() 永久打开闸门。如果想重新等待下一轮需要 clear() 重置标志但这轮已经 wait 的线程不会重新阻塞需要新建或处理边界条件。4.2 用 Semaphore 限制同时执行的线程数量Semaphore 维护一个计数器。每次 acquire 减一成功进入每次 release 加一。当计数器降到 0acquire 会阻塞直到其他线程 release。典型场景是限制并发数量。比如下游接口只能承受 3 个并发请求就用初始值为 3 的信号量控制。import threading import time sem threading.BoundedSemaphore(3) def limited_task(name): with sem: print(f{name} 开始执行, flushTrue) time.sleep(1) print(f{name} 执行结束, flushTrue) threads [threading.Thread(targetlimited_task, args(ftask-{i},)) for i in range(6)] for t in threads: t.start() for t in threads: t.join()虽然这里启动了 6 个线程但同一时刻最多只有 3 个线程能进入 with sem 代码块其余线程在 Signal 处排队。普通 Semaphore 允许 release 次数超过初始值会把计数器越加越高。BoundedSemaphore 多了一个边界检查release 超过初始值会抛 ValueError。在工程中更推荐 BoundedSemaphore。4.3 daemon 线程与主线程退出策略创建 Thread 时可以设置 daemon 参数。daemon 为 True 的线程会在主线程退出时被强制终止不会等待它执行完。这在后台轮询、健康检查、日志异步上报场景中很有用。import threading import time def background_loop(): while True: print(后台任务运行中, flushTrue) time.sleep(1) t threading.Thread(targetbackground_loop, daemonTrue) t.start() time.sleep(3) print(主线程退出)主线程 sleep 3 秒后就会退出daemonTrue 的后台线程随之停止。如果这里不设置 daemonTrue即使主线程逻辑结束程序也会一直等待后台线程无法退出。需要区分的是daemon 线程并不是“后台任务”的代名词。它意味着这个线程随时可能被主线程结束打断不能在 daemon 线程里执行必须落盘或回滚的事务逻辑。像耗时但要求一致性的任务应该用普通线程加 join 等待完成。4.4 通过线程枚举信息定位卡住问题threading.active_count 返回当前存活线程数量threading.enumerate 返回所有存活线程列表threading.current_thread 返回当前调用线程对象。正常退出前可以打印这些信息判断是否有线程没有结束。import threading import time def job(): time.sleep(2) t threading.Thread(targetjob, namelong-task) t.start() time.sleep(0.5) print(存活线程数:, threading.active_count()) for thread in threading.enumerate(): print(thread.name, thread.daemon) t.join()如果程序卡住优先把这段信息放在疑似阻塞点之前观察哪些线程还活着哪些线程已经退出。线程名、daemon 状态、是否 join都是排查的关键信息。5. 多线程实战用固定数量的 worker 并发处理一批任务5.1 需求设计与为什么要先拆消费模型假设有一批资源需要处理例如读取一批本地文件、检查一批配置、调用一批允许访问的数据接口。逐个处理太慢于是决定用多线程并发。最常见也最稳的模型是生产者消费者模型。主线程把任务放进 queue多个 worker 线程从 queue 取任务并执行。这样有几个好处。线程数量和任务总数解耦。任务很多时不会盲目创建几千个线程而是固定 5 个或 10 个 worker避免线程切换开销。worker 之间天然通过 queue 隔离不需要额外保护任务列表。任务不断加入时worker 仍然可以自动处理程序结构不会变化。先明确边界再写代码任务是什么worker 是什么任务正常完成怎么标记任务抛出异常怎么兜底主线程什么时候确认全部完成。5.2 基础版本queue.Queue 与 worker 循环下面的代码实现了“12 个任务3 个 worker”的处理流程。为了模拟 IO 等待每个任务休眠 0.5 秒如果换成真实网络请求或文件读取只需要替换 do_work 内部代码。import queue import threading import time task_queue queue.Queue() def do_work(worker_name): while True: task task_queue.get() if task is None: break try: print(f{worker_name} 开始处理 {task}, flushTrue) time.sleep(0.5) print(f{worker_name} 完成 {task}, flushTrue) finally: task_queue.task_done() def start_workers(worker_count): workers [] for idx in range(worker_count): t threading.Thread(targetdo_work, args(fworker-{idx},)) t.start() workers.append(t) return workers def stop_workers(workers): for _ in workers: task_queue.put(None) for t in workers: t.join() def main(): for i in range(1, 13): task_queue.put(f资源-{i}) workers start_workers(3) start_time time.time() task_queue.join() stop_workers(workers) elapsed time.time() - start_time print(f全部任务执行完成耗时 {elapsed:.2f}s) if __name__ __main__: main()关键点分为三层。worker 为什么是 while True。因为线程启动后要持续等待队列中的新任务不能只处理一个任务就退出。退出条件是通过哨兵 None 实现的每个 worker 消费一个 None 后 break这样能保证所有 worker 都退出。先调用 task_queue.join() 是等待队列里所有任务都被执行完毕join 返回后再发 None否则如果任务还没消费完就发哨兵有可能出现 worker 提前退出、剩余任务无人处理的情况。task_done 必须放在 finally 中。即使任务执行时抛异常也要通知队列这个任务已经结束否则主线程的 join 会一直等待。do_work 里 try/finally 的组合就是为异常兜底设计的。5.3 扩展结果收集与异常兜底上面例子只打印线程执行过程。实际业务通常需要把每个任务的结果收回来。可以引入一个共享 list并用 Lock 保护追加操作。为了保证异常日志可见还需要在 worker 里捕获 Exception。import logging import queue import threading import time logging.basicConfig(levellogging.INFO) task_queue queue.Queue() result_list [] result_lock threading.Lock() def do_work(worker_name): while True: task task_queue.get() if task is None: break try: # 这里替换成真实 IO 任务 time.sleep(0.5) result f{task} 处理结果 except Exception: logging.exception(%s 处理 %s 失败, worker_name, task) result f{task} 处理失败 finally: task_queue.task_done() with result_lock: result_list.append((worker_name, task, result))收集结果时不只是在 finally 之后 append。Lock 保护的是 result_list 这个共享对象避免多个线程同时修改变量时互相覆盖。异常捕获要尽量精确不要捕获后什么都不记录。即使业务允许失败也要在日志里留下任务标识和异常堆栈。如果任务顺序很重要不能在结果收集时直接依赖完成顺序。方案有两种给每个任务放一个任务编号worker 处理时把编号一起写回或者是完成后统一按编号排序。5.4 运行验证看结果和耗时运行基础版本代码时预期结果是12 个任务3 个 worker每个任务 sleep 0.5 秒。理论串行耗时 6 秒左右三个 worker 并发后耗时接近 2 秒到 2.5 秒。控制台会看到 worker-0、worker-1、worker-2 交错处理不同资源输出顺序不固定。这是正常现象不代表代码有问题。验证并发是否真正生效可以先按单串行方式跑一遍记录耗时再用相同任务量启动 3 个 worker对比总耗时。结果验证不能只看程序是否正常结束。至少还要确认所有任务都进入了 worker没有任务残留在队列中。结果列表数量等于任务总数。日志里没有未捕获异常的堆栈。多次执行结果保持一致说明没有明显竞争问题。5.5 真实项目还要补哪些约束上面的示例可以跑通但离生产环境还有一定距离。真实项目里worker 通常不会直接替换 time.sleep而是要处理更复杂的 IO 调用。生产环境要注意以下差异关注点学习脚本生产脚本线程数量随意设置任务少根据下游连接池、API 限制、CPU 情况设置任务队列长度不限制使用 maxsize 限制队列长度防止内存膨胀异常处理捕获打印记录结构化日志、失败任务单独排队或重试超
返回列表