ARTICLE DETAIL

资讯详情

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

Python多进程并发编程:IPC、进程池与现场排查

Python多进程并发编程:IPC、进程池与现场排查 搞网络并发编程的朋友迟早都会碰到“进程”这根硬骨头。不管是爬虫批量抓数据、写推送网关还是把慢接口从主服务里拆出来异步执行你都会发现单线程顶不住磁盘和网络互相等着线程又容易踩 Python 的 GIL 限制于是“多进程”就成了最直观、也最稳定的并发手段。这篇是《Python网络并发编程》系列的第二篇专门讲进程从 multiprocessing 最基础的 Process 和 Pool 讲起一路聊到网络服务器里多进程 accept 的经典模型、进程间通信IPC、僵尸进程和端口占用这些现场问题。适合已经能用 Python 写基本网络程序的读者看完之后至少能回答三件事什么场景该用进程、进程之间怎么传数据、进程挂了/卡了/不见了应该怎么排查。1. 先搞明白进程在网络并发里到底扮演什么角色1.1 进程不是一个“更重的线程”那么简单进程是操作系统分配资源的最小单位线程是 CPU 调度的最小单位这句话背下来容易真正理解它的分量得从“隔离”说起。每个进程都有独立的虚拟内存空间、独立的文件描述符表、独立的全局变量和栈。一个进程里某个第三方 C 扩展直接段错误只会把自己干崩其他进程纹丝不动线程则共享同一进程的地址空间一个线程把指针写到野地址上整个进程都可能一起陪葬。在网络场景下隔离的意义比单机计算更明显。我以前写爬虫就吃过亏某次用多线程抓全网商品数据一个不起眼的解析库在解析畸形 HTML 时把解释器干崩了所有线程全部中断几百个请求的上下文全丢。后来改成“每个 worker 一个进程”单个 worker 崩了主控程序依然活着还能标记这个任务失败、重试、写日志。这个差异单机跑 demo 感觉不到一上生产马上就体现出来。另一个不得不提的事实是进程在多核 CPU 上是真的“并行”而愿意牺牲地址空间隔离换来的线程在这个层面并不占优。很多初学者以为 Python 多线程能直接吃满多核结果发现 8 个线程跑 CPU 密集任务CPU 使用率只有 100% 左右就是在围着 GIL 打转。进程则不同每个 Python 子进程都有自己独立的解释器实例自带一个 GIL所以 4 个进程跑计算任务只要机器核数足够CPU 使用率就能接近 400%。这就是为什么“进程”在网络并发编程里从来不是备选项而是绕不开的核心手段。1.2 GIL 到底卡在哪多进程怎么绕过去的GIL 全称 Global Interpreter Lock全局解释器锁。在 CPython 里同一时刻只允许一个线程执行 Python 字节码所以纯计算逻辑的多线程代码实际表现和单线程相当还额外背上了线程切换开销。不是没有人尝试移除 GIL但 CPython 为了保住 C 扩展兼容性和 JIT 生态至今仍默认保留。好在 GIL 只锁解释器不锁操作系统资源如果你在做文件读写、网络 IO、sleep 等待这些操作会把 GIL 释放掉所以 IO 密集型任务用多线程仍然很香。明白这个底层设计之后再说多进程思路就清晰了既然解释器实例之间不共享 GIL那我就直接开多个解释器实例。Python 提供了multiprocessing模块底层有两种实现路径一种是fork()复制当前进程的内存快照另一种是spawn重新导入主模块并启动一个新的解释器。子进程之间通过 pickle 序列化传递参数通过管道、队列、共享内存交换数据。代价是什么呢每次启动一个子进程都要重新初始化解释器环境内存占用几十 MB 起步IPC 还有序列化和反序列化成本。所以“多进程”不是免费的你的每一点并行优势都得用资源换回来。1.3 进程、线程、异步三者的分工别搞对立网络并发里最容易被问懵的问题就是那我到底该用进程、线程还是 asyncio我的习惯是先把任务分成两种特性吃不吃 CPU以及等不等 IO。CPU 密集计算、压缩、加解密、图像处理、回测指标计算这类任务在 Python 里优先考虑进程。IO 密集HTTP 请求、数据库读写、文件流、长连接推送因为线程和 asyncio 都擅长等待选谁看并发规模和业务复杂度。高连接数边缘网关、聊天服务这种动不动几万连接的场景用线程容易栈内存爆掉用进程更不现实事件循环asyncio几乎是唯一合理选项。真实业务往往不是单一种类。比如爬虫网络请求是 IO 密集解析 HTML 又带点 CPU 密集还要隔离第三方解析库的崩溃风险。于是最稳妥的架构是“主控进程跑调度 进程池做解析 连接层用异步/线程池控制 IO”。把这些模型组合起来比争论“哪个更好”有意义得多。我在后面的实战里也会按这个思路展开。2. 上手 multiprocessing从 Process 到 Pool2.1 最小的多进程程序Process 与 joinimport multiprocessing import os def worker(num: int): pid os.getpid() print(f子进程 {pid} 处理任务 {num}) if __name__ __main__: ctx multiprocessing.get_context(fork) procs [ctx.Process(targetworker, args(i,)) for i in range(4)] for p in procs: p.start() for p in procs: p.join() print(主进程结束)几个容易被新手忽略的点join()的作用是让主进程阻塞等待子进程退出。没join的话子进程可能还在运行主进程就先退出了如果子进程是daemonTrue主进程退出时它会被强杀。if __name__ __main__一定要写。Windows 和 macOS 上multiprocessing默认使用spawn方式启动子进程它会重新执行主模块如果没有这层保护子进程又去启动子进程会无限递归最后报 RuntimeError。为什么我显式写了get_context(fork)Linux 上默认就是 fork不写也没问题但当你需要代码跨平台时最好明确启动方式。fork启动最快因为它直接复制父进程内存缺点是把线程锁、socket 状态等一并继承容易出隐性 bugspawn干净但每次启动要重新导入模块慢forkserver用一个后台服务进程专门负责创建子进程兼顾速度和干净但在 Windows 上不可用。写网络服务时我个人倾向于用spawn保证状态干净再配合进程池减少启动次数。2.2 任务结果怎么拿Queue 和 Manager进程最大的问题是函数return出来的结果主进程是拿不到的。子进程和主进程的地址空间独立return只能回到子进程自己的世界里。所以要么用队列、管道、共享内存要么直接交给ProcessPoolExecutor这种高层封装。先用 Queue 版本import multiprocessing as mp def worker(q, num): q.put(num * num) if __name__ __main__: q mp.Queue() procs [mp.Process(targetworker, args(q, i)) for i in range(4)] for p in procs: p.start() for p in procs: p.join() results [q.get() for _ in range(4)] print(results)这个版本能跑但有个坑如果某个worker抛了异常它不会向队列放数据而主进程还在傻等q.get()于是程序永远阻塞。更稳的做法是用concurrent.futures.ProcessPoolExecutorfrom concurrent.futures import ProcessPoolExecutor, as_completed def heavy(x: int) - int: return x * x with ProcessPoolExecutor(max_workers4) as pool: futures [pool.submit(heavy, i) for i in range(8)] for fut in as_completed(futures): result fut.result() # 子进程异常会在这里重新抛出 print(result)ProcessPoolExecutor会在主进程侧捕获子进程的异常fut.result()抛出的异常和普通调用几乎一致方便日志和告警。我个人写新项目时很少直接操作裸Process去收集结果进程池优先。Manager 的用法m mp.Manager() shared_list m.list() shared_dict m.dict()它启动了一个独立的 manager 服务进程其他进程通过代理访问它管理的对象。适合低频的数据交换比如汇报每个 worker 的进度。代价是每次访问都有一层代理开销比共享内存慢不少别拿它存高频计数器。2.3 进程池参数怎么定一个靠谱的参考公式进程池max_workers到底填多少没有银弹但我可以给你一套经过项目验证的决策流程。先判断任务类型。CPU 密集任务worker 数建议等于 CPU 核数或cpu_count() 1。多出来的那一个是为了应对某些进程短暂进入系统调用、让出 CPU 时带来的调度间隙。实测在大多数 Linux 机器上cpu_count() 1比严格等于核数时吞吐更好看。IO 密集任务可以用2 * cpu_count() 1这个经典方案起步。这个公式来自并发社区的经验不是一定最优而是给你一个不会错的起点然后通过监控 CPU 空闲率、队列积压数量去调整。爬虫这类场景最终瓶颈往往在目标网站的并发限制和本机端口资源上盲目加 worker 只会换来一堆超时。内存是硬约束。每个 Python 进程的固有开销就算很精简也有 20MB 到 30MB加上 pickle 缓冲、任务数据100 个进程轻松吃掉 3GB 内存。所以我设计服务时进程数的选择顺序永远是先看内存预算再看 CPU 核数最后看外部依赖的并发上限。任务类型推荐进程数注意事项CPU 密集cpu_count()或cpu_count()1不要超过物理核太多否则上下文切换吃光收益IO 密集数据库/磁盘从2*cpu_count()1开始监控数据库连接数和磁盘队列深度网络请求爬虫/API从 CPU 核数左右起步逐步加压注意目标站点频率限制、本机端口上限混合型分层设置网络层异步/线程 计算层进程池3. 进程通信到底怎么选IPC 四种姿势IPC 这个词看着高大上其实就一句话让两个拥有独立地址空间的进程交换数据。热词榜上那“electron 主渲染进程 IPC”说的是 JS 桌面应用的前后台通信本质上和咱 Python 里讲的队列、管道是同一层抽象。理解了一组概念别的技术栈只是换了个接口而已。3.1 Queue / JoinableQueue最像日常工作的通信方式multiprocessing.Queue的底层是一个 pipe 加一个锁外加一个 feeder 线程。数据从 put 进去之后先被 pickle 序列化然后写进管道消费者在另一端反序列化后取出。import multiprocessing as mp def producer(q, n): for i in range(n): q.put(i) q.put(None) # 生产结束信号 def consumer(q): while True: item q.get() if item is None: break print(f消费 {item}) if __name__ __main__: q mp.Queue(maxsize100) p1 mp.Process(targetproducer, args(q, 10)) p2 mp.Process(targetconsumer, args(q,)) p1.start(); p2.start() p1.join(); p2.join()常见的坑有三个队列里的对象必须能被 pickle 序列化。lambda、嵌套函数、生成器对象统统不行。大对象序列化开销明显。你往队列里塞一个 10MB 的字符串整体耗时和内存都会翻倍还不如让它留在共享内存或磁盘临时文件里。JoinableQueue多一个task_done()和join()它能帮你确认队列里的每条消息都已经被消费者处理完了而不是仅仅被取走。这在“等待所有任务完成后再做汇总”的场景里很好用。3.2 Pipe低延迟双向通道Pipe返回一对连接对象适合两个进程一对一通信。duplexTrue表示双向都能发False表示一个只读一个只写。parent_conn, child_conn mp.Pipe()父进程parent_conn.send(...)子进程child_conn.recv(...)反过来也可以这取决于duplex参数。如果你只需要单向把duplexFalse能省掉一些内部锁。Pipe 最大的优势是简单、低延迟小消息的性能比 Queue 好。最大的坑是管道缓冲区有限如果一端狂发、另一端不读发送端会阻塞如果两端都在互相等待对方读取就死锁了。我见过一个同事写的 demo父子进程各发一个大字典双方都先 send 再 recv结果是两边都卡在 send 上。解决办法很简单先约定好消息顺序让一端只发、另一端只收或者限制单条消息大小。3.3 Value/Array Lock共享内存里的原子操作如果只是数数、改个状态位用队列就觉得重管道又嫌麻烦那就上共享内存。multiprocessing.Value和Array直接分配一块内存多个进程可以读写同一块数据。import multiprocessing as mp def add(lock, counter, n): for _ in range(n): with lock: counter.value 1 if __name__ __main__: counter mp.Value(i, 0) lock mp.Lock() procs [mp.Process(targetadd, args(lock, counter, 10000)) for _ in range(4)] for p in procs: p.start() for p in procs: p.join() print(counter.value) # 期望 40000这里i表示 C 的 int 类型d表示 double类型码和 Python 的array模块一致。没有锁的情况下counter.value 1不是原子操作读取、加一、写回三步之间可能被另一个进程插一脚最后结果会小于 40000。这个 bug 很阴间因为它大多数时候是对的偶尔少几个数排查起来特别费劲。共享内存的另一个用途是给 worker 传只读的常量大数组。比如量化回测里所有进程都要读同一份行情数据与其每个进程各存一份不如用Array放一份进程只去读不需要锁。数据量小的话直接放args里随进程复制过去更省心。3.4 Manager共享对象的瑞士军刀Manager()会启动一个 Server 进程代理它管理的所有共享对象其他进程通过代理访问。它能管的不只是 list 和 dict连Namespace、Lock、Queue都能管。我自己的使用场景是“进度上报”进程池里每个 worker 把自己的状态写进m.dict()主控进程每隔几秒读一次在控制台上画个进度条。对这种低频小数据交换Manager 写起来代码最短也没有序列化限制的烦恼。缺点是慢因为每次读写都要穿越进程边界、走代理协议。你要是拿它做高频计数器性能会很难看。低频元数据、共享结构、跨平台兜底——这是 Manager 的正确打开方式。3.5 一张表总结四种 IPC 怎么选通信方式适用场景性能复杂度主要坑Queue多生产者/多消费者任务分发中等带序列化低pickle 限制、消息堆积Pipe一对一双向通信高低延迟低缓冲区满死锁Value/Array Lock计数器、共享结构、大数组高近原生低必须加锁类型码Manager进度共享、复杂结构低代理调用中性能慢启动开销大不管你选哪种心里都要有个数进程通信传输的一定是“数据”不是“对象”。对象在发送端被序列化成字节流在接收端被还原成新对象两者只是值相等绝非同一个东西。想通了这一点就不会写出“传一个连接对象过去”这种不切实际的代码。4. 网络服务里的多进程模型实战4.1 fork 之后子进程 accept经典的预派生模型网络服务里的多进程最常见的目标就是“多进程同时 accept 同一个监听 socket”。这在 Linux 的 C 语言网络编程里有很经典的 pre-fork 模型父进程 bind listen fork 一批子进程子进程各自accept()等待连接。TCP 协议栈和 socket 在内核里保证了多个进程同时 accept 不会把同一个连接分给两个进程内核会自动做负载均衡唤醒最近 sleep 的进程。Python 里用multiprocessing也能实现类似结构import socket import multiprocessing as mp def worker(srv: socket.socket, stop_event: mp.Event): srv.settimeout(1.0) while not stop_event.is_set(): try: conn, addr srv.accept() except socket.timeout: continue except OSError: break with conn: data conn.recv(1024) conn.sendall(data.upper()) if __name__ __main__: srv socket.socket(socket.AF_INET, socket.SOCK_STREAM) srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) srv.bind((0.0.0.0, 8080)) srv.listen(64) stop_event mp.Event() procs [mp.Process(targetworker, args(srv, stop_event), daemonFalse) for _ in range(4)] for p in procs: p.start() try: while True: stop_event.wait(1) except KeyboardInterrupt: stop_event.set() for p in procs: p.join(timeout3) srv.close()这段代码在 Linux 上可以用子进程继承了监听 socket。几个细节说一下子进程里srv.settimeout(1.0)很重要没有它accept 会永久阻塞主进程设置 stop_event 也没法让子进程醒来退出。except OSError: break用于监听 socket 被关闭或错误时退出循环。daemonFalse配合显式事件退出比直接把子进程设成 daemon 然后靠父进程退出杀死更可控。4.2 SO_REUSEADDR 与 SO_REUSEPORT绑同端口的两种思路网络服务重启时老服务还没跑完的 TCP 连接会留下大量 TIME_WAIT 状态连接。TIME_WAIT 是 TCP 协议主动关闭方在等待 2MSL最大报文段生存时间后释放资源的必要状态Linux 上默认大约 60 秒。如果服务端自己也参与主动关闭就会有 TIME_WAIT 留在那个端口上此时再 bind 同端口会报Address already in use。设置SO_REUSEADDR就是为了让 TIME_WAIT 状态下的端口允许重新绑定这是每个 TCP 服务端都该加的选项。SO_REUSEPORT则更进一步它允许多个进程各自 bind 同一个 IP端口内核在收到新连接时做负载均衡。这样每个进程都有自己独立的监听 socket不需要靠 fork 继承进程退出也不会影响别的进程。代码里用法if hasattr(socket, SO_REUSEPORT): srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1)注意 SO_REUSEPORT 在 Linux 和 macOS 上可用Windows 上支持有限写跨平台服务时要判断。如果你用 SO_REUSEPORT各个进程可以独立启动比如通过进程管理器管理多个实例每个实例自己 bind 同端口。这比 fork 的优点是进程之间解耦更强缺点是需要内核支持且所有进程必须都设置 SO_REUSEPORT 才能同时绑定成功。4.3 子进程退出后的“收尸”问题僵尸进程与 SIGCHLD多进程服务运行久了最容易遇到的就是僵尸进程。正常情况下子进程退出后父进程需要调用wait()/waitpid()回收它的退出状态这个调用会释放进程表中残留的记录。如果父进程一直不调用退出后的子进程就会处于defunct/Z状态也就是僵尸进程。僵尸进程不能被kill -9杀死因为它已经死了只留了个登记簿在系统里。如果僵尸进程大量堆积进程表被占满新的进程就没法创建了。pre-fork 模型里如果某个子进程被外部信号杀掉而父进程只在那里循环 wait不管子进程僵尸就会产生。最常见的补救是注册 SIGCHLD 信号处理器import signal, os def reap(signum, frame): while True: try: pid, status os.waitpid(-1, os.WNOHANG) except ChildProcessError: break if pid 0: break signal.signal(signal.SIGCHLD, reap)-1表示等待任意子进程os.WNOHANG表示如果没有子进程退出就立即返回这样循环可以把所有退出子进程都收干净。这个 handler 要在创建子进程之前注册父进程收到 SIGCHLD 后一有子进程退出就回收僵尸。4.4 进程数到底开多少要结合网络模型看在 pre-fork 模型里子进程数不是一个固定公式能解决的我通常按这个思路来定如果请求处理是纯 CPU 计算子进程数接近 CPU 核数即可。如果请求处理涉及大量外部 IO查 Redis、调第三方接口、读写磁盘可以适当增加但你要监控每个 worker 的 IO 等待。超过一定数量后瓶颈会变成外部服务的连接池限制或自身的文件描述符额度。文件描述符上限也很关键。每个 TCP 连接至少消耗一个 fd系统默认ulimit -n可能是 1024这是很多测试环境里“连接数一高就全挂”的元凶。上线前记得调大并检查ulimit -SHn 65535是否在部署脚本里生效。我自己的经验预派生 4 个 worker 处理 100 并发请求绰绰有余因为每个请求大多是几毫秒的 IO 等待4 个进程轮流照顾好几百个连接没问题。真正到了百万连接级别就不是进程数问题了而是需要引入事件循环模型让每个进程内部再去管理成千上万个连接。5. 现场排查进程相关的常见事故笔记热词榜里那些“查询8080端口进程”“WPS进程无法关闭”“微信运行好多进程呀”“msmpeng 后台进程过大”的问题往深了看其实都和“进程生命周期管理”有关。这一节我不只讲理论而是把处理过的几个现场问题原原本本写上。5.1 端口被占用到底是谁占了 8080先来最常被问的问题启动服务时报Address already in use怎么排查按系统不同命令也不一样Linuxss -lntp sport :8080或lsof -i :8080能直接看到 PID 和进程名。macOSlsof -i :8080默认可能没有ss用 lsof 就行。Windowsnetstat -ano | findstr :8080最后一列是 PID然后用任务管理器翻对应进程。找到 PID 之后先确认这个进程是不是上一次启动的残留如果是自己的服务没退干净看看它的父进程是什么尽量通过接口或信号优雅关闭。如果确实查不出有用信息可以看它的启动时间对比一下是不是在最近一次部署前后出现的。不要一言不合就kill -9尤其当这个进程可能是别的服务托管的你杀了它会触发自动重启看起来就是“怎么都杀不死”。我写过一个小工具函数在服务启动前主动探测端口占用并打印 PID 和完整命令行能省去很多部署时候的来回沟通import socket, subprocess def check_port_in_use(port: int): with socket.socket() as s: try: s.bind((0.0.0.0, port)) return False except OSError: out subprocess.check_output( [lsof, -i, f:{port}, -sTCP:LISTEN, -n, -P], textTrue, ) print(out) return True5.2 僵尸进程看不见的进程表杀手排查僵尸进程的命令ps -eo pid,ppid,stat,comm | grep -E Z|defunct看到 STAT 是 Z或者 COMMAND 是defunct就说明有僵尸。僵尸的 PPID 指向谁谁就是没有及时回收的父进程。修复方法也是围绕父进程展开要么父进程里加 waitpid 循环要么重启父进程让孤儿被 init 收养init 会处理它们的回收。在 Python 里我通常还是用前面提过的 SIGCHLD handler并顺手在监控脚本里对Z状态进程做告警因为进程表的容量是有限的僵尸堆积会让新的Process/fork失败。5.3 子进程卡住、收不到数据这类问题我列一个快速排查清单子进程是不是没退出导致join()一直等先看ps -ef | grep 你的脚本名。队列里还有没有数据如果消费者提前结束put端可能因为管道写满而阻塞。消息是不是太大序列化/反序列化耗时超过网络 IO导致任务实际吞吐远低于预期。是不是死锁比如两个进程各自持有锁又互相等对方释放。如果子进程 CPU 正常、也不退出用py-spy dump --pid pid直接看 Python 调用栈能快速揪出卡在哪个库的哪一行。 Linux 上还可以用strace -p pid看系统调用卡在 read/write 还是 poll判断它是在等待队列、管道还是网络。5.4 别把“进程多”当病毒多进程架构的正常与异常热词里“微信运行好多进程呀”其实一点不稀奇。现代软件喜欢把一个整体服务拆成多个进程比如主界面进程、渲染进程、网络进程、更新进程好处是单点崩溃不会拖垮全部也方便系统按进程级别分配资源。Windows 里的msmpeng是系统自带的防病毒进程它偶尔占用高有它自己的逻辑“WPS 进程无法关闭”常常是因为它的托盘进程、模块进程在设计上会互相拉起你杀了一个它又起来一个。这些现象放到我们自己的 Python 服务里对应的教训是多进程程序要设计明确的退出协议不能让子进程被父进程kill -9后变成孤儿然后继续占着端口、连接池。不要试图在业务代码里用“杀掉进程树”这种粗暴办法解决问题它大概率会在你的数据文件写到一半时留下半截状态。真正要做的是区分后台服务进程、业务 worker、守护进程的职责谁崩溃谁来拉起谁来回收退出状态上线前用脚本把所有可能的退出码列出来测试一遍。6. 再进阶一点进程 异步的组合6.1 事件循环 进程池网络并发的高吞吐形态很多服务单用多进程会发现一个问题进程里的每个 worker 还是同步阻塞地处理请求如果请求里有大量外部 IOworker 就又闲又等。这时候把异步和多进程叠加起来效果会好很多——不是非此即彼而是各管一段。一个很典型的结构是主进程/主线程跑 asyncio 事件循环负责成千上万个连接的事件分发遇到真正吃 CPU 的计算任务比如 JSON 大报文解析、模板渲染、加解密通过进程池交出去执行事件循环则继续服务其他请求。import asyncio from concurrent.futures import ProcessPoolExecutor def heavy_parse(data: bytes) - dict: import json return json.loads(data) async def handle(reader, writer, pool): data await reader.read() loop asyncio.get_running_loop() result await loop.run_in_executor(pool, heavy_parse, data) writer.write(str(result).encode()) await writer.drain() writer.close() async def main(): pool ProcessPoolExecutor(max_workers4) try: server await asyncio.start_server( lambda r, w: handle(r, w, pool), 127.0.0.1, 8888, ) async with server: await server.serve_forever() finally: pool.shutdown() if __name__ __main__: asyncio.run(main())这里run_in_executor(pool, ...)会把函数体交给进程池asyncio 本身不会阻塞。进程池里的计算是同步的但主循环还在继续 accept 其他连接。这个组合在 Python 社区非常常见任务队列系统里的 worker、异步框架里的高级接口大多都是这套思路。6.2 多进程结果的收敛别让每个进程都写数据库多进程跑任务最直接的结果归集方式是每个 worker 处理完了就把结果写到数据库。看着简单实际容易踩坑数据库连接数上限。一个 worker 建一个连接50 个 worker 就是 50 个连接还没算应用自身的连接池数据库很容易被打爆。写入打爆热键。所有 worker 同时写一张表锁竞争激烈写入延迟飙升。事务边界混乱。一个任务失败数据可能写了一半没有统一回滚逻辑。更合理的结构是引入“结果收敛层”worker 把结果塞到消息队列 / Redis list / 本地文件一个专门的结果写入进程批量消费攒够一批再执行批量 INSERT。批量写往往能比单条写提升数倍速度而且能统一做去重、校验和兜底重试。量化回测这种场景更是如此多进程回测不同参数组合每个进程算出一堆收益曲线和指标千万别让每个回测进程直接往库里写几千行曲线数据而是每个进程把曲线成批发回主进程主进程汇聚后统一落库。这既让主进程能看到整体进度也方便做参数对比和可视化。6.3 什么情况下不要用进程写到最后必须泼一盆冷水进程不是万能的很多时候它甚至是负优化。任务太轻太短。一个任务就是做一次字符串拼接进程启动、调度、IPC 序列化的开销已经比任务本身还大。这种场景用 asyncio 或线程就够了。你已经用 asyncio 写了纯异步程序只是没吃满 CPU那加进程池只因为你想用多核。记得只把真正重的计算扔进程而不是把整个 loop 塞进去否则就是拿着异步代码去办同步的事。机器内存太小开不了几个进程就 OOM这时候要么上异步要么优化单进程资源占用。我见过一个线上事故有人给一个批量接口加了个 16 进程的池每个进程都加载一份几百 MB 的词表8GB 内存直接被打满服务 OOM。教训是共享大对象尽量用Array或落到Redis让所有进程只读同一份数据而不是各复制一份。7. 我的一点实操心得文章快写完了按老规矩收个尾说点真正让我少交学费的习惯。第一我写任何和进程相关的代码都会在启动入口先写一份“进程生命周期清单”哪个进程负责创建、哪个负责回收、子进程退出条件是什么、主进程收到 SIGTERM/SIGINT 后谁先退。这份清单不写进文档而是直接抽象成一个统一的管理函数统一处理 start、stop_event、join、reap。后来线上所有任务都走这个入口僵尸进程和端口残留发生的次数基本降到了零。第二监控优先于排查。没有监控进程跑了多少、内存多高、僵尸几个都是黑的。我至少会在服务里加一套基础指标当前 worker 数、每进程 RSS、accept 队列积压用 py-spy 定期采样或者直接让子进程上报到 Manager dict。等真的出问题时这些指标就是事故现场的第一手证据。第三上线之前一定会用strace或py-spy跑一轮冒烟测试。别嫌土很多所谓的“莫名其妙卡死”在strace下面十秒就现原形——不是卡在阻塞的accept()就是卡在没设超时的queue.get()。把这些工具放进你的工具箱比记住十篇理论文章都管用。下一讲我打算聊网络并发里的 asyncio 和协程把“多进程 异步”怎么组合出高吞吐服务拆开讲细。到时候你会发现进程是骨架异步是肌肉两者配合才是 Python 网络并发最实用的姿势。
返回列表