多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Python并发编程实战:线程队列线程池与生产者消费者模式详解

Python并发编程实战:线程队列线程池与生产者消费者模式详解 一直用Python写脚本处理数据的人大概率都会遇到一个场景一个任务明明可以拆成多个独立的小任务同时跑结果却只能排队一个个执行眼看着CPU闲着、接口在干等、耗时翻了好几倍心里那个急。这时候多数人第一反应是“多线程并行跑”但动手之后又发现线程之间数据一竞争就乱套动不动还死锁甚至写了半天也不知道该用线程还是进程。这篇就围绕线程、队列、生产者与消费者、线程池这几个Python并发编程里的核心主题把怎么组织代码、怎么避免踩坑、怎么让任务真正高效跑起来讲清楚。适合刚接触并发编程的初学者也适合写了一段时间但总感觉哪里不对、想系统理一遍的开发。线程和队列的组合本质上是良药线程负责干活队列负责传话。把两者结合起来再上一个线程池基本就能覆盖90%以上的日常多任务处理需求。接下来我按自己的实际使用经验把这套东西从底层原理到落地写法一步步拆开说。1. Python线程基础与GIL1.1 线程到底是什么很多初学者会把线程理解成“性能加速器”觉得加线程代码就跑得快这个认知如果不矫正后面会踩很多坑。线程本质上是一个进程内部的执行单元一个进程可以包含多个线程它们共享这个进程的内存空间包括全局变量、堆数据、打开的文件句柄等。用生活化的类比来说进程像一个公司线程就是公司里的员工。公司有独立的办公场地内存空间员工之间共用会议室、公共资料区共享数据但每个员工有自己的工牌和茶杯栈空间与局部变量。线程最大的优势是轻量创建开销比进程小得多切换也快所以适合处理那些需要同时等待多个外部资源的场景。Python标准库里的threading模块提供了线程的基本封装一个最简单的线程写法如下import threading import time def worker(name): for i in range(3): print(f线程 {name} 正在工作: {i}) time.sleep(0.5) t1 threading.Thread(targetworker, args(A,)) t2 threading.Thread(targetworker, args(B,)) t1.start() t2.start() t1.join() t2.join() print(所有线程执行完毕)这里start()会启动线程join()会阻塞主线程直到对应线程结束。需要注意join()是保证主线程不会提前退出导致后续逻辑错乱的关键操作尤其是当线程还没有跑完、主线程就要收集结果或者继续做其他事时少了join()往往会出现诡异时序问题。1.2 GILPython线程绕不开的坎聊Python多线程GIL全局解释器锁是绕不开的话题。GIL是CPython解释器中的一个互斥锁它保证同一时刻只有一个线程能够执行Python字节码。也就是说多个线程在单核层面轮流抢CPU时间片本质上并不能把计算任务真正并行化。网上很多说法很夸张说“Python多线程是废物”这个判断过于绝对。真实情况是如果你的任务是CPU密集型的比如大量数值计算、循环处理多线程不仅不会提速还可能因为锁竞争和上下文切换变得更慢这种情况应该用multiprocessing多进程。但如果任务是I/O密集型的比如网络请求、数据库读写、文件读写线程在等待I/O时会释放GIL其他线程能继续执行多线程就有非常明显的收益。一个常见的实践参考假设你要发起100个HTTP请求每个请求耗时1秒串行需要100秒。用线程池开10个线程并发请求总耗时大约能压缩到10秒出头。原因就是网络请求的等待时间里GIL被释放了其他线程可以跑。但如果换成用Python做10万次大数乘法多线程可能比串行还慢。所以第一步思维转变是面对任务先分清它是CPU密集型还是I/O密集型再决定用线程还是进程。1.3 线程安全问题因为线程共享进程内存所以多个线程同时读写同一个变量时会出现竞态条件。举个最常见的例子import threading counter 0 def increment(): global counter for _ in range(1000000): counter 1 threads [threading.Thread(targetincrement) for _ in range(5)] for t in threads: t.start() for t in threads: t.join() print(f最终值: {counter}期望值: 5000000)这里的counter 1看似一行其实底层至少包含三步读取当前值、加1、写回。多个线程同时执行这三步时可能同时读到一个相同的旧值再加1写回导致丢更新最终结果小于5000000。解决这个问题的基础手段是加锁Python提供threading.Lock用with上下文管理器锁住临界区。但锁用多了又可能引发死锁这也是为什么我更推荐优先用下面要讲的队列来传递数据而不是让多个线程直接操作共享变量。队列内部自带了锁和同步机制在多线程场景下是天然安全的数据交换通道能少写很多Lock代码。2. 队列线程安全的通信桥梁2.1 queue.Queue的基本用法queue模块是Python标准库中专门为多线程通信设计的线程安全队列。它的好处在于你不需要自己加锁put和get方法内部已经处理了同步问题。最基础的使用方式如下import queue q queue.Queue(maxsize10) # 放入数据 q.put(任务A) q.put(任务B) # 取出数据 item q.get() print(item) # 任务A这里要注意maxsize参数它控制队列最多能放多少个元素。如果队列满了put会阻塞直到有空位如果队列空了get会阻塞直到有数据进来。这种天然的阻塞特性恰好是生产者消费者模式的基础。很多时候我们还会配合put_nowait和get_nowait使用它们不会阻塞如果操作无法完成会抛出queue.Full或queue.Empty异常适合需要非阻塞轮询的场景。我平时用得最熟练的三件套是task_done()和join()的组合。q.join()会阻塞主线程直到队列中的所有元素都被get取走并且每个元素都调用了task_done()。这在主线程需要等待队列任务全部完成的场景中非常关键。比如一批任务放进队列消费者从队列取走并处理主线程用q.join()等所有任务结束比用time.sleep()硬等靠谱得多。2.2 队列的三种形态queue模块其实提供了三种不同排序规则的队列类型排序规则适用场景queue.Queue先进先出FIFO最常用任务按提交顺序公平处理queue.LifoQueue后进先出LIFO类似栈最新任务优先适合需要优先处理新任务的场景queue.PriorityQueue优先级排序小值优先每个元素是(优先级, 数据)元组适合根据重要程度处理任务PriorityQueue实际使用时有个小坑如果两个任务的优先级相同后续元素会继续比较元组本身内容如果数据不可比较会抛TypeError。稳妥的做法是给每个任务带上序号降低比较冲突的概率类似(优先级, 序号, 数据)。2.3 队列的底层是怎么回事从原理上说queue.Queue内部使用了threading.Condition来实现条件变量的同步。Condition本质上是一个锁加一个等待通知机制当队列为空时get会进入等待状态让出CPU和锁当有数据被放进来时put会执行notify唤醒一个等待的消费者。这种方式比单纯的忙碌轮询高效得多不会空转消耗CPU。简单理解就是消费者不再每隔几毫秒来看一次锅里有没有饭而是锅里的饭好了之后自动接到一个“开饭通知”。关于队列的容量选择我个人的经验是maxsize不要设置得太小否则生产者会被阻塞得过于频繁导致整体吞吐下降但也不要设置得太大否则一旦消费者崩溃任务会大量堆积在内存里。日常爬虫或任务分发场景队列设为任务总数的1.5到2倍就已经比较合理了。如果是无限流直接不传maxsize但必须自己监控队列长度防止内存被吃满。3. 生产者与消费者模式3.1 模式的价值在于解耦与削峰生产者消费者模式是并发编程里最经典也最实用的架构模式。它把“产生数据的一方”和“处理数据的一方”通过一个缓冲队列隔离开。生产者的职责是往队列里放数据消费者的职责是从队列里取数据并处理二者之间不需要知道对方的任何细节也不需要同步等待对方。这个模式的价值体现在三个层面第一是解耦生产者不用关心谁消费、怎么消费消费者也不用关心数据从哪来第二是缓冲如果生产速度和消费速度不一致队列可以暂存数据避免一方被另一方拖慢第三是削峰瞬时涌入大量任务时队列能平滑处理压力避免瞬间把所有线程打满。用现实例子来理解就像餐厅后厨和前厅之间有一个传菜口。厨师做菜不需要站在前台等客人吃完客人点完餐也不需要一直等厨师传菜口把两边的节奏拉开谁快谁慢都不互相影响。很多秒杀系统、日志收集系统、任务调度系统核心架构就是这个模型。3.2 一个具体可跑的示例下面用一个多生产者多消费者示例来演示场景是模拟爬虫下载100个页面import queue import threading import time import random q queue.Queue(maxsize20) def producer(worker_id, count): for i in range(count): task_id ftask-{worker_id}-{i} q.put(task_id) # 模拟产生任务的时间 time.sleep(random.uniform(0.01, 0.05)) print(f生产者 {worker_id} 完成) def consumer(worker_id): while True: task q.get() if task is None: # 哨兵信号 q.task_done() break # 模拟处理任务的时间 time.sleep(random.uniform(0.1, 0.3)) print(f消费者 {worker_id} 处理完成: {task}) q.task_done() print(f消费者 {worker_id} 退出) producers [threading.Thread(targetproducer, args(i, 25)) for i in range(4)] consumers [threading.Thread(targetconsumer, args(i,)) for i in range(5)] for t in producers: t.start() for t in consumers: t.start() # 等待所有生产者完成 for t in producers: t.join() # 发送停止信号 for _ in consumers: q.put(None) # 等待消费者退出 for t in consumers: t.join() print(所有任务处理完毕)这里用了一个很实用的技巧哨兵信号None。因为消费者是无限循环执行get的如果队列空了它们会一直阻塞等待主线程就没法优雅退出。往队列里放和消费者数量相同个数的None每个消费者取到None就退出循环干净利落。需要注意这里的q.task_done()在哨兵分支也要调用一次因为q.join()会等待所有取出的元素都执行task_done()如果漏了哨兵的task_done主线程用q.join()等待时就会永久阻塞。3.3 生产消费速度不匹配怎么办实际运行时常常出现一种情况生产者太快队列容量很快就满了生产者被阻塞在put上反过来消费者太快队列空了消费者阻塞在get上。这两种情况都是正常的不必强行去改。真正需要关注的是生产者和消费者的线程数量如何配比。经验上如果生产一个任务耗时远小于处理一个任务耗时比如生产耗时0.01秒、处理耗时0.2秒那消费者的数量基本要往生产速度除以处理速度的量级去配。上面例子里生成4个生产者、5个消费者是因为处理时间大约比生产时间多5到10倍。如果是I/O密集型处理线程数还可以适当放大我实际用下来Python线程在I/O等待时能很好利用GIL释放机制消费者的数量开到CPU核心数的3到5倍都不算过分但也要结合网络、数据库连接数等外部资源上限来定。4. 线程池频创线程不如复用线程4.1 为什么需要线程池如果按照上面手动threading.Thread一个个创建线程任务一多很快就会暴露问题线程创建和销毁本身有系统开销线程数量失控会让操作系统调度压力变大内存占用也会上升而且手动管理线程生命周期很容易出错比如线程忘记join就提前退出或者线程数量无限增长导致系统资源耗尽。线程池的核心思想是复用预先创建一批线程让它们持续等待任务有任务就执行没有任务就挂着不需要反复创建销毁。Python 3.2之后concurrent.futures模块提供了ThreadPoolExecutor这是官方推荐的线程池实现底层封装了队列、线程管理、任务调度这些细节。我实际项目中基本都用这个很少再手动裸写threading。4.2 ThreadPoolExecutor基本用法下面是最常用的提交方式submit会返回一个Future对象代表一个异步执行的任务from concurrent.futures import ThreadPoolExecutor, as_completed import time def fetch(url): # 模拟网络请求 time.sleep(1) return f数据来自 {url} urls [fhttps://example.com/page/{i} for i in range(20)] with ThreadPoolExecutor(max_workers8) as executor: futures {executor.submit(fetch, url): url for url in urls} for future in as_completed(futures): url futures[future] try: data future.result() print(f完成 {url} - {data}) except Exception as e: print(f任务 {url} 失败: {e})这里有两个关键点。第一with语句会在代码块结束时自动调用线程池的shutdown(waitTrue)等待所有提交的任务执行完毕再释放线程资源比手动管理线程优雅得多。第二as_completed会在有任务完成时立即返回对应的Future不会因为某个任务特别慢而阻塞其他已完成结果的获取。4.3 submit与map的取舍除了submitThreadPoolExecutor还提供了map方法用法类似内置的map函数但是是并发执行with ThreadPoolExecutor(max_workers4) as executor: results list(executor.map(fetch, urls))两者的主要区别是map会按传入顺序返回结果整体像同步执行第一个任务没跑完你在遍历结果时就得等submit配合as_completed则可以在任意任务完成时就立刻拿到结果适合任务耗时差异大、希望尽快处理结果的场景。如果只是批量执行并且需要保证结果顺序和输入顺序一致用map更简洁如果需要逐个处理异常、或者任务之间耗时差异很大用submit更灵活。4.4 max_workers到底设置多少这是新手问得最多的问题也是线程池最容易拍脑袋配的参数。经验上分两种情况如果是I/O密集型任务线程数可以设置成任务并发目标数量但也要受到外部资源瓶颈约束。比如目标是同时发起50个HTTP请求max_workers50是可行的但如果是操作数据库连接池连接上限是10那线程数最多也就10否则线程都在排队等待连接。如果拿不准外部资源上限可以先从CPU核心数 * 5开始跑并发测试再根据实际耗时调整。如果是CPU密集型任务Python线程受GIL限制线程数设置成CPU核心数甚至更低反而更好想真正并行计算就得走ProcessPoolExecutor那是另一个话题。ThreadPoolExecutor内部维护了一个任务队列和一个工作线程集合提交给线程池的任务会先进入队列再由空闲线程取出执行。max_workers控制了同时能跑的任务数量队列可积压的任务数量则没有硬性上限。所以如果任务提交速度远大于消费速度任务会在内存中积压这一点在后面问题排查部分还要专门说。5. 清单实战爬虫场景完整示例5.1 综合需求拆解把前面讲的线程、队列、生产者消费者、线程池整合起来最典型的场景就是爬虫或者批量接口请求。我拿一个实际的清单来演示整体架构需要抓取500个商品详情页页面URL由一个基础列表接口分批生成抓取页面需要访问详情接口抓完要把结果写入文件。如果串行做网络等待时间会让整体耗时非常感人。用线程池加队列改造之后代码结构会更清晰扩展性也好。5.2 完整代码import queue import threading import time import json import random from concurrent.futures import ThreadPoolExecutor q queue.Queue(maxsize200) def fetch_detail(task): url task[url] # 模拟网络请求耗时 time.sleep(random.uniform(0.5, 1.5)) detail { id: task[id], url: url, title: f商品标题 {task[id]}, price: random.randint(50, 500) } return detail def produce_tasks(): for i in range(500): task {id: i, url: fhttps://example.com/goods/{i}} q.put(task) print(任务生产完毕) def consume_task(): while True: task q.get() if task is None: q.task_done() break try: detail fetch_detail(task) save_lock.acquire() try: with open(goods.json, a, encodingutf-8) as f: f.write(json.dumps(detail, ensure_asciiFalse) \n) finally: save_lock.release() except Exception as e: print(f处理失败 {task}: {e}) finally: q.task_done() print(消费者工作线程退出) save_lock threading.Lock() producer_thread threading.Thread(targetproduce_tasks) producer_thread.start() with ThreadPoolExecutor(max_workers10) as executor: futures [executor.submit(consume_task) for _ in range(10)] for f in futures: f.result() # 等生产者线程结束后再停止 producer_thread.join() print(全部任务执行完毕)这段代码有几个细节值得多说一句。一是文件写入必须加锁因为多个线程同时写一个文件很容易出现内容交错或者写入失败这里用save_lock保护了open和write这一段临界区。二是消费者本质上也是用线程池管理的每个工作线程都在无限循环中从队列取任务直到拿到None才退出这样线程池的生命周期是可控的。三是q.task_done()写在finally里保证即使任务处理抛异常q.join()也不会被卡住这个习惯非常重要。5.3 架构上还能怎么扩展上面的代码已经能直接跑了但在真实项目中生产者也不会只从一个地方来。常见扩展思路包括生产者从多个API分页拉数据那就是多个生产者线程消费者处理的内容需要写多个库表可以把队列分成多条子队列每个消费者组处理一条子队列如果系统还要持久化消息、保证重启后不丢任务那就要引入Redis Stream或专业消息队列来做更可靠的任务存储。从生产者消费者模式的角度理解这些扩展其实都是在队列两侧增加“生产者来源”和“消费者去向”核心的“队列解耦”思路始终不变。6. 常见问题与排查技巧实录6.1 死锁线程死锁大概是并发编程里最烦的问题。典型场景是多个线程互相持有对方需要的资源不释放比如线程A持有锁1等待锁2线程B持有锁2等待锁1两个线程都卡死了。排查死锁最直接的办法是抓取线程转储。在Python里可以用faulthandler或py-spy不过日常简单的排查可以先检查自己的锁获取顺序。一条非常有效的铁律所有线程获取多个锁的顺序保持一致。比如约定必须先获取锁1再获取锁2所有线程都按这个顺序来就能从逻辑上避开循环等待。另外能用队列就不用锁队列内部已经处理好了线程安全这也是我前面一直强调队列好用的原因之一。6.2 任务堆积任务堆积常见于生产者生产速度大于消费者消费速度线程池内部队列被填满内存逐渐吃紧。排查时可以定期打印队列大小或线程池工作状态import threading import time def monitor(q): while True: print(f队列剩余: {q.qsize()}) time.sleep(2) t threading.Thread(targetmonitor, args(q,), daemonTrue) t.start()q.qsize()能拿到当前队列长度配合时间戳上报到监控系统就可以看到任务积压趋势。解决任务堆积一般有三种思路一是增加消费者线程数量二是增大队列容量、削峰填谷更从容三是优化消费者的处理逻辑把同步处理改成异步批量处理比如攒够N条结果再统一写库。6.3 线程池任务异常被吞很多人在使用ThreadPoolExecutor时遇到一个很坑的现象任务函数内部抛了异常程序竟然没有报错看起来好像任务没有执行。原因是Future的异常不会自动抛出来只有调用future.result()时异常才会被重新抛出。如果只submit不处理Future异常就被静默吞掉了。解决方法是遍历as_completed拿取每个Future并显式调用result()或者给任务函数内部加try/except做兜底。我自己习惯在任务函数最外层包一个统一的异常捕获函数把异常信息打日志既保证了任务不中断又方便事后排查。6.4 join与task_done使用时机混淆join这个命名在threading.Thread.join()和queue.Queue.join()里含义完全不同初学者最容易搞混。Thread.join()是等某个线程结束Queue.join()是等队列中的任务全部被处理完。两者经常配合使用但必须理解各自的语义。使用Queue.join()时如果忘记对get取出的任务调用task_done()队列会永远认为任务没有完成join()就会永久阻塞。这个问题非常隐蔽因为程序不会报错看起来就是停在那里不动。我排查过多次这类问题最后都是发现某个分支没有写task_done()。建议在所有get之后的代码路径里无论成功还是失败都在finally里调用task_done()。6.5 线程安全的外部资源线程池里如果共享了同一个HTTP会话、同一个数据库连接、同一个日志文件句柄这些外部资源本身是否线程安全必须确认。比如requests.Session在并发环境下并不是绝对安全的一般建议每个线程一个Session或者加锁复用。数据库连接池也是同理连接对象通过池化管理不能多个线程直接共用一个连接。这块没有统一解法核心思路是线程之间共享的对象如果是不可变的那最安全如果是可变对象要么加锁要么每个线程单独创建副本。7. 我的一些使用心得动手写并发之前先想清楚任务类型I/O密集就用线程CPU密集就用进程不确定就先做个小实验用time.time()跑一段串行对比并发看真实数据再决定不要让网上那些“Python多线程没用”的论调替你做了决定。从简单的ThreadPoolExecutor开始用别一上来就整复杂的任务框架。生产者消费者模式加队列基本能解决绝大多数实际业务问题而且结构清晰、容易调试。等到出现真正的性能瓶颈再考虑用ProcessPoolExecutor弥补CPU密集型短板或者用asyncio改写协程方案到时候你已经有清晰的线程模型底子理解起来会快很多。整个线程、队列加生产者消费者的体系别看聊起来好像概念多真正落地也就是几行代码的事。抓住两个关键一是队列是线程之间传数据的最安全通道能少写很多锁二是线程池负责管理线程生命周期能让你专注业务逻辑不用天天操心线程创建销毁。把这两个点吃透并发编程这块基本就不会出大乱子。
返回列表