尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Python多进程并行执行组件:让CPU密集型任务跑满多核

Python多进程并行执行组件:让CPU密集型任务跑满多核 做批量任务的朋友应该都对“CPU打不满、时间线性增长”这件事有切肤之痛。我这次分享的并行执行组件进程版就是为了把大批量任务拆分到多个进程里同时跑而写的核心包括进程池管理、任务队列分发、结果回收、异常自动拉起源码也放在文末提到的完整工程里了。做数据清洗、量化回测、爬虫抓取、文件转码这类场景的朋友直接拿去改一改就能用。一句话说明白它解决什么问题单进程跑任务8核CPU只有1个核在忙其他人围观用这个组件任务分发到多个子进程每个核都有活干整体耗时基本能做到原来的几分之一。本文会从选型思考、架构设计、关键代码、踩坑记录到实测调参把进程级并行这件事讲透。1. 为什么是进程——把任务模型理清楚再动手很多人一上来就纠结“线程好还是进程好”“协程是不是更轻”其实答案完全取决于任务类型。我在设计这个组件之前先花了一小时把任务模型盘清楚这里把判断逻辑直接给大家。1.1 线程和协程绕不开的那堵墙先说线程。线程最大的问题是“安全”和“抢占”多线程共享同一块内存空间一旦有共享变量没加锁数据错乱只是时间问题。调试线程问题非常折磨人崩溃的现场往往和实际原因隔了十几行代码打印出来的变量值已经是串位的。协程则是“单线程内想办法”。协程的切换确实轻但本质还是在事件循环里排队对于大规模CPU密集计算帮助非常有限它擅长的是I/O等待场景比如大量网络请求时把时间片让给别的协程去发下一个请求。进程走的是完全不同的路线每个子进程有独立的地址空间变量、对象都是副本主进程怎么改都影响不到子进程从根上避免锁竞争。进程由操作系统内核调度可以真正分散到不同CPU核心上跑。单进程崩溃不影响其他人主进程只需要重启它就好。代价就是进程创建和销毁比线程重一些进程间通信没有共享内存那么直接。这就是进程池存在的意义——把建进程的昂贵开销摊到很多次任务上。1.2 哪类任务天生就该用进程拆分根据我的经验下面几类任务特别适合进程化CPU密集型计算比如批量图像处理、特征提取、矩阵运算、文件编码转换。单线程跑就是“1核干活、7核围观”进程化后基本能线性提速。隔离要求高的任务跑某些不稳定的第三方库崩溃一次可能会带崩整个主程序。放到子进程里崩了拉起一个就行主进程稳如泰山。内存占用大的任务子进程独立地址空间跑完释放不会在主进程里留下莫名其妙的引用链。重量级初始化后重复执行的任务比如加载一个几百MB的模型加载一次要好几秒如果每个任务都重新加载就浪费了。进程池可以先让子进程初始化好环境然后用队列不断丢任务进去执行。判断标准很简单任务越长、越密集越值得用进程。如果单个任务执行只要几毫秒那进程创建和通信的耗时反而会淹没收益这类短平快的任务用线程池或协程更合适。2. 组件架构与代码结构拆解确定用进程方案后我先画了一张逻辑上的数据流图然后按照职责拆分模块。这里先给整体架构后面再逐段解析源码。2.1 核心模块划分完整的组件划分成四个模块模块职责关键点任务队列从主线程接收待执行任务按序分发使用跨进程队列任务对象必须可序列化进程池管理维护一组常驻子进程控制生命周期控制最大并发数回收空闲进程任务调度与回收分配任务到空闲进程收集结果和异常通过结果队列回传支持超时控制守护与自愈监控子进程状态异常退出自动重启处理僵尸进程、维护最小工作进程数为什么要把“进程池管理”和“任务调度与回收”拆成两个模块因为它们的关注点完全不同进程池只关心“现在有几个可用进程、谁空闲、谁忙碌”任务调度则关心“任务给谁、结果去哪个队列”。混在一起代码会长成一个难以维护的大泥球排查问题时也不好定位。2.2 一次完整并行执行的生命周期从外部看这个组件的用法很简洁executor ProcessParallelExecutor(max_workers4) executor.start() futures executor.submit_all(task_fn, task_list) results executor.wait_all(futures, timeout30) executor.shutdown()流程可以拆成这几步调用方调用submit_all把所有任务打包放入发送队列。每个空闲子进程从发送队列取一个任务执行task_fn拿到结果后放入结果队列。主进程后台线程从结果队列回收结果挂到对应Future对象上。wait_all阻塞等待所有Future完成或者超时返回已经完成的部分。调用方处理后调用shutdown主进程发送退出信号等待所有子进程退出。注意第4步的超时设计这是实际业务中特别实用的功能。有个视频转码场景个别视频文件损坏导致单条任务卡死没有超时控制整个组件都被拖住了。加了超时后超时的任务标记为失败并返回其余正常任务的结果照常取回不会被一颗老鼠屎坏一锅汤。3. 核心实现与源码解析接下来是重点我照着实际代码逐块解析。考虑到组件要同时兼容Linux和Windows我在底层直接用multiprocessing库实现没有依赖concurrent.futures因为ProcessPoolExecutor在任务粒度、取消控制、异常细节上不如自己写的灵活。3.1 进程池骨架与任务队列设计进程池的核心是保持固定数量的子进程常驻避免每次任务都重复创建进程。创建进程的操作很昂贵——Linux下要复制整个进程地址空间哪怕用了copy-on-write也要走一遍内核的fork流程Windows下更夸张等于重新初始化一次解释器。import multiprocessing import queue import threading import time from concurrent.futures import Future class ProcessParallelExecutor: def __init__(self, max_workersNone, task_queue_size1000): self.max_workers max_workers or multiprocessing.cpu_count() self.task_queue multiprocessing.Queue(maxsizetask_queue_size) self.result_queue multiprocessing.Queue() self._workers [] self._futures_map {} self._result_collector_thread None self._running False这里有几个设计取舍要说清楚。task_queue设置maxsize是为了防止任务生产速度远超消费速度导致队列无限膨胀吃掉所有内存。如果任务源来自数据库读取一万个任务排队没执行积压的内存可能先让主进程OOM。result_queue不设上限因为结果回收线程会持续取数据理论上不太会积压真遇到结果生产者比消费者快很多的情况也可以给它加个下限控制。future映射表_futures_map是任务ID到Future对象的映射。为什么不用列表因为子进程在执行任务前会先从任务队列取到任务并把一个唯一ID放回结果队列里这样主进程回收结果时才知道“这个结果是哪个任务的”。用字典做映射可以实现O(1)的查找一万个任务也不会慢。3.2 任务提交、结果回收与超时控制的完整实现任务提交那一步我给每个任务生成唯一ID封装成一个元组(task_id, task_name, task_args, task_kwargs)塞进任务队列。同时给Future对象登记到映射表并把任务ID绑定在Future的属性上。def submit_all(self, task_fn, tasks): futures [] for task in tasks: task_id uuid.uuid4().hex fut Future() fut.task_id task_id self._futures_map[task_id] fut self.task_queue.put((task_id, task_fn.__name__, task)) futures.append(fut) return futures这里我特意把task_fn本身没有放进队列只放了它的名字。因为跨进程传递函数对象很麻烦pickle对函数序列化有天然限制只有模块级函数能正常picklelambda和嵌套函数直接解析就报错。我做的约定是客户端对任务函数统一命名子进程里再按文件名或注册表找到对应函数执行。def _worker_loop(self, worker_id, task_runner): self._log(fworker {worker_id} started) while True: try: task_id, fn_name, task self.task_queue.get(timeout1) except queue.Empty: continue try: result task_runner(task) self.result_queue.put((task_id, success, result)) except Exception as exp: self.result_queue.put((task_id, error, repr(exp)))结果回收在主进程单独起了一个守护线程它不断从结果队列读取消息根据消息里的task_id找到对应Future再调用future.set_result或future.set_exception。这里最重要的一点Future的set_result必须在主进程的线程里调用不能直接在回收线程里调用因为这会触发Future绑定的回调函数回调可能需要与主线程交互。我用一个专用的回收线程来做这个事天然隔离。超时控制是基于Future自带的future.result(timeoutN)实现的。wait_all只是把所有Future包进concurrent.futures.wait传入总超时时间然后返回完成与未完成的Future集合def wait_all(self, futures, timeoutNone): from concurrent.futures import wait, FIRST_COMPLETED done, pending wait(futures, timeouttimeout, return_whenALL_COMPLETED) return done, pendingwait内部的实现会在所有future完成或超时后返回不会提前返回。需要部分结果时可以用return_whenFIRST_COMPLETED循环拉取。3.3 子进程守护、崩溃自愈与优雅退出子进程跑着跑着挂了是常态内存不够、任务函数触发段错误、第三方C扩展库崩了都会让子进程直接消失。刚开始我没做自愈结果跑一天后进程池空了一半没人管。后来补上了心跳监控。实现思路是每个任务开始前子进程先向结果队列发一条心跳消息内容是(worker_id, heartbeat, timestamp)。主进程维护一张worker_id - 最近心跳时间的表一个后台线程每5秒扫描一次发现某个worker超过阈值比如60秒没心跳就标记该worker失效然后启动一个新worker补位。def _heartbeat_monitor(self, max_idle60): while self._running: now time.time() dead_workers [] for worker_id, last_beat in self._heartbeats.items(): if now - last_beat max_idle: dead_workers.append(worker_id) for worker_id in dead_workers: self._restart_worker(worker_id) time.sleep(5)这个方案属于“轻量级自愈”照顾了绝大多数场景。如果Worker长期没有任务执行导致心跳停止那就必须先判断是“真死”还是“假死”。真死是进程退出了我们还可以通过process.is_alive()判断假死是任务卡死比如死循环、IO阻塞心跳监控本身并不能强制杀掉卡死的任务。遇到这种情况我建议在任务队列层面的每个任务上再套一个timeout参数子进程内部使用multiprocessing.Queue的get(timeout...)来中断卡死的任务函数——但这个前提是任务函数本身能被timeout打断纯CPU死循环是打断不了的这类任务得靠上层业务自己设断点。优雅退出同样有很多细节。直接调process.terminate()虽然暴力有效但会丢失子进程内未输出的日志和未写回的结果只用process.join()不传超时的话任务卡死的子进程会一直不退出主进程跟着卡死。我的做法是双阶段关闭def shutdown(self, timeout10): self._running False for _ in range(self.max_workers): self.task_queue.put(None) # 哨兵值通知子进程退出 for worker in self._workers: worker.join(timeouttimeout) # 给子进程优雅退出的时间 for worker in self._workers: if worker.is_alive(): worker.terminate() # 兜底强制结束 self._result_collector_thread.join(timeout5)重点在于哨兵值None。子进程的worker_loop每次get时先判断拿到的值是不是None是就主动退出循环然后进程自然结束。任务卡死的子进程收不到正常退出信号会在join(timeout)之后被terminate()兜底。注意terminate()之前最好留一点时间让子进程完成手头的工作否则正在写入的文件可能只写了一半。4. 从日志里挖出来的坑框架写好了只是开始真正运行起来的坑才让人头大。下面这几个问题都是我真实遇到过、排查过甚至深夜加班解决过的每一行都是血泪。4.1 僵尸进程与“内存只增不减”的真相Linux下的僵尸进程特别隐蔽。子进程正常退出后会短暂变成僵尸态Zombie如果父进程不及时wait()回收僵尸态会一直保留占用进程表项。进程表被占满后系统就创建不了新进程了。我在这个组件里已经有join()回收机制理论上不会漏但曾经在一条异常分支里漏调用了join导致跑了一天后系统里堆了几十个僵尸进程。排查方法一言难尽ps -ef | grep defunct发现一堆defunct标记再查他们的PPID都是我的主进程。修复的方式是在所有子进程结束路径上都要有join兜底包括异常分支。我现在直接在shutdown()里对每个worker写一遍join(timeout)再配合terminate()确保子进程都被回收。内存量只增不减的问题则来自队列。任务积压在task_queue里队列里的Python对象占内存进程中结果对象也存在_futures_map里任务跑完忘记从映射表删掉内存就不会释放。这个组件跑大数据集时会明显感觉到内存随着任务增多而膨胀。解决办法两个任务完成后立刻从_futures_map里删除映射队列设置合适的maxsize生产速度过快时阻塞提交方。4.2 子进程里打日志导致“进程一起卡死”这是模板级经典坑。子进程把日志打到sys.stdout而stdout底层缓冲区是跨进程共享的多个子进程同时写就会在一个锁上死等。日志多的时候几个子进程全部卡在写日志上任务队列里堆积越来越多父进程也感知不到子进程假死表现就是“任务一个都不完成了”。解决方式是把子进程的日志重定向到指定文件或者使用logging模块并给每个子进程配置独立的日志文件。我后来做了个统一的日志管理器按worker_id命名日志文件既避免了输出竞争排查具体任务时还能按进程单看日志效率反而提升。4.3 任务分发不均与“长尾效应”任务大小差异大时最简单的一次全分发策略会让快的进程闲下来慢的进程还在跑最后几个大任务整体等待时间被长尾任务拖长。跑某个图像处理任务时绝大多数图处理只要0.1秒极少数超大图要跑5秒最后同步等待时间全花在最大图上。改进方法是改成动态拉取模式任务不预先分发给子进程而是放在共享队列里每个子进程完成当前任务后自己去取下一个。这样快的进程自动多干几个慢的进程少干几个整体完成时间明显缩短。这个组件最终采用的就是动态拉取策略效果在前面已经看到——整体耗时从几十秒降到十几秒。4.4 为什么关闭后还有残留子进程这个坑出现在Windows上主进程shutdown()之后任务管理器里仍有子进程残存。原因是子进程的daemon标志没设好。在multiprocessing中只有设置daemonTrue的子进程才会在主进程退出时被强制终止否则Windows上主进程退出时子进程会继续默默活着。还有一个隐蔽因素是子进程内部又派生了孙进程孙进程不受daemon约束即使父进程终止它自己还会继续跑。后续遇到这类残留根本解决办法是任务函数里不要在自己的进程内部再用multiprocessing开子进程或者把子进程的逻辑统一收敛到顶层worker中。5. 实测数据与调参建议最后放一组实测数据以及我压测之后沉淀下来的参数经验方便大家拿到组件后直接做初步调参。5.1 单核基线对比测试结果测试任务选择了纯CPU密集计算的哈希循环任务模拟数据指纹计算10万条数据每条计算量相当。环境为8核机器、Linux系统Python 3.10。执行方式总耗时平均CPU利用率说明单进程串行186s~12%单核打满其余核围观线程池8线程184s~15%受GIL影响几乎无提升进程版组件4 worker51s~45%4核参与计算进程版组件8 worker27s~85%8核基本都用上了用GIL特性的线程池几乎等于没救进程才是这类任务的正解。CPU密集场景下进程数从4增加到8我的实测加速接近线性186秒缩到27秒快了大约6.9倍。没有到理论上的8倍是因为进程调度、队列通信、结果Pickle传输都有固有开销。5.2 worker数量与队列大小的建议worker数量不是越多越好。我的经验公式是CPU密集任务worker数取CPU核心数I/O密集或包含大量睡眠等待的任务worker数取核心数的2到4倍。注意我这句话依赖一个隐蔽前提CPU密集任务在子进程里不会频繁调用线程池等机制抢GIL。如果任务函数内部本身又用了多线程那就要按“子进程内再核数”重新算不能只看外部worker数。Windows上cpu_count()常会返回逻辑核心数如果机器有超线程可以先用os.cpu_count() // 2做降级启动。队列大小建议默认100到1000之间取决于单个任务的数据量。每个任务是几百字节的小对象时队列可以设大一点每个任务是几MB的数组时队列设大会直接把内存吃爆这种情况建议50以内让任务边取边填利用流水线优势降内存峰值。最后分享几点我的实际体会这套组件从最初一百来行的demo发展到后来带自愈、超时、日志隔离、优雅关闭的“完整版”中间踩了很多坑也让我对进程调度和multiprocessing底层的理解深了不少。个人最深的体会是先判断任务类型再选并发模型是最值钱的一步方案错了后面代码写再好都白搭装饰器、观察者模式的业务封装尽量后置第一版先把进程池的骨架跑通再叠加业务语义。还有一个很管用的小技巧给子进程传入的唯一ID我除了用它匹配Future还会把它写进每条日志的前缀。排查问题的时候同一任务的日志聚在一起串行流程一目了然。早期的日志混杂在一起满屏飞消息处理进度根本对不上号这个细节建议读者自己实现时也顺手加上。后续如果要扩展可以按几个方向走任务拆分为有依赖关系的有向无环图再实现执行调度把任务队列换成持久化中间件实现断点续跑给子进程增加内存和CPU资源限制防止单条野任务拖垮整台机器。不管怎么扩展进程池框架只要打得足够干净换存储、换调度策略都不至于伤筋动骨。
返回列表