☰
Python多进程并发编程:IPC、进程池与现场排查
2026/9/30 4:37:20 网站建设 项目流程

搞网络并发编程的朋友,迟早都会碰到“进程”这根硬骨头。不管是爬虫批量抓数据、写推送网关,还是把慢接口从主服务里拆出来异步执行,你都会发现,单线程顶不住磁盘和网络互相等着,线程又容易踩 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 到 Pool

2.1 最小的多进程程序:Process 与 join

import 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(target=worker, args=(i,)) for i in range(4)] for p in procs: p.start() for p in procs: p.join() print("主进程结束")

几个容易被新手忽略的点:

  • join()的作用是让主进程阻塞等待子进程退出。没join的话,子进程可能还在运行,主进程就先退出了;如果子进程是daemon=True,主进程退出时它会被强杀。
  • if __name__ == "__main__"一定要写。Windows 和 macOS 上multiprocessing默认使用spawn方式启动子进程,它会重新执行主模块;如果没有这层保护,子进程又去启动子进程,会无限递归,最后报 RuntimeError。
  • 为什么我显式写了get_context("fork"):Linux 上默认就是 fork,不写也没问题;但当你需要代码跨平台时,最好明确启动方式。fork启动最快,因为它直接复制父进程内存,缺点是把线程锁、socket 状态等一并继承,容易出隐性 bug;spawn干净但每次启动要重新导入模块,慢;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(target=worker, 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.ProcessPoolExecutor:

from concurrent.futures import ProcessPoolExecutor, as_completed def heavy(x: int) -> int: return x * x with ProcessPoolExecutor(max_workers=4) 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(maxsize=100) p1 = mp.Process(target=producer, args=(q, 10)) p2 = mp.Process(target=consumer, args=(q,)) p1.start(); p2.start() p1.join(); p2.join()

常见的坑有三个:

  • 队列里的对象必须能被 pickle 序列化。lambda、嵌套函数、生成器对象统统不行。
  • 大对象序列化开销明显。你往队列里塞一个 10MB 的字符串,整体耗时和内存都会翻倍,还不如让它留在共享内存或磁盘临时文件里。
  • JoinableQueue多一个task_done()和join(),它能帮你确认队列里的每条消息都已经被消费者处理完了,而不是仅仅被取走。这在“等待所有任务完成后再做汇总”的场景里很好用。

3.2 Pipe:低延迟双向通道

Pipe返回一对连接对象,适合两个进程一对一通信。duplex=True表示双向都能发,False表示一个只读一个只写。

parent_conn, child_conn = mp.Pipe()

父进程parent_conn.send(...),子进程child_conn.recv(...),反过来也可以,这取决于duplex参数。如果你只需要单向,把duplex=False,能省掉一些内部锁。

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(target=add, 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(target=worker, args=(srv, stop_event), daemon=False) 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(timeout=3) srv.close()

这段代码在 Linux 上可以用,子进程继承了监听 socket。几个细节说一下:

  • 子进程里srv.settimeout(1.0)很重要,没有它,accept 会永久阻塞,主进程设置 stop_event 也没法让子进程醒来退出。
  • except OSError: break用于监听 socket 被关闭或错误时退出循环。
  • daemon=False配合显式事件退出,比直接把子进程设成 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,怎么排查?按系统不同,命令也不一样:

  • Linux:ss -lntp sport = :8080或lsof -i :8080,能直接看到 PID 和进程名。
  • macOS:lsof -i :8080,默认可能没有ss,用 lsof 就行。
  • Windows:netstat -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"], text=True, ) print(out) return True

5.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 子进程卡住、收不到数据

这类问题我列一个快速排查清单:

  1. 子进程是不是没退出,导致join()一直等?先看ps -ef | grep 你的脚本名。
  2. 队列里还有没有数据?如果消费者提前结束,put端可能因为管道写满而阻塞。
  3. 消息是不是太大?序列化/反序列化耗时超过网络 IO,导致任务实际吞吐远低于预期。
  4. 是不是死锁?比如两个进程各自持有锁,又互相等对方释放。
  5. 如果子进程 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 还是同步阻塞地处理请求,如果请求里有大量外部 IO,worker 就又闲又等。这时候把异步和多进程叠加起来,效果会好很多——不是非此即彼,而是各管一段。

一个很典型的结构是:主进程/主线程跑 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_workers=4) 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 网络并发最实用的姿势。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询