前面介绍了多线程编程,由于 Python 中,CPython 存在 GIL(全局解释器锁),同一个进程内,同一时刻仅有一个线程执行 Python 字节码。CPU 密集型场景,多线程无法利用多核并行,提速效果极差,甚至比单线程更慢。
为了解决 CPU 密集任务并发问题,推出了多进程编程。多进程可以创建独立的子进程,每个进程有自己独立的内存空间,完全并行,不受 GIL(全局解释器锁)限制。
os.fork() 是 Unix/Linux/macOS 系统上创建进程的底层方法,Windows 不支持。
调用 os.fork() 瞬间,操作系统复制当前父进程,生成全新子进程。子进程会复制父进程的代码段、数据、堆、文件描述符、环境变量。父进程、子进程拥有独立内存空间,修改变量互不影响。
执行 fork() 后,两个进程从 fork() 调用处同时向下执行。
通过 fork() 的返回值区分父子:
返回 0:当前是子进程
返回 大于 0 的正整数:当前是父进程,数字是子进程 PID
返回 负数:创建进程失败
例如:
import os
print("主进程 PID:", os.getpid())
pid = os.fork()
if pid == 0:
print(f"我是子进程,PID={os.getpid()},父进程 PID={os.getppid()}")
else:
print(f"我是父进程,PID={os.getpid()},创建了子进程 PID={pid}")如果在 Windows 中运行上述代码,输出如下:
主进程 PID: 23236
Traceback (most recent call last):
File "d:\share_dir\workspace\5.demo\python_demo\demo.py", line 5, in <module>
pid = os.fork()
^^^^^^^
AttributeError: module 'os' has no attribute 'fork'如果在 Linux 中运行上述代码,输出如下:
snow@loclahost:~/python$ python3 demo.py
主进程 PID: 1487
我是父进程,PID=1487,创建了子进程 PID=1488
我是子进程,PID=1488,父进程 PID=1487注意,fork() 是底层操作系统方法、与系统强相关、不跨平台,实际开发不推荐。不使用 fork(),我们该如何创建进程,实现多进程编程,推荐使用 multiprocessing 模块,底部的细节由模块取处理。
multiprocessing 模块是 Python 官方跨平台多进程标准库,功能强大,使用方式和 threading 几乎一样。
下面介绍两种创建进程的方式:
直接将一个函数传递给 Process(),和线程创建方式一样,例如:
from multiprocessing import Process
import time
def task(name):
print(f"子进程 {name} 运行")
time.sleep(2)
print(f"子进程 {name} 结束")
if __name__ == "__main__":
# 创建一个进程
p1 = Process(target=task, args=("A",))
# 创建另一个进程
p2 = Process(target=task, args=("B",))
p1.start()
p2.start()
p1.join() # 等待完成
p2.join() # 等待完成
print("所有进程执行完毕")运行代码,输出如下:
子进程 B 运行
子进程 B 结束
子进程 A 运行
子进程 A 结束
所有进程执行完毕注意:上面代码是跨平台的,Windows、Linux 执行效果一样。
创建一个自己的类,继承 Process 类,实现进程创建,例如:
from multiprocessing import Process
import time
# 继承 Process 类
class MyProcess(Process):
def __init__(self, name):
super().__init__()
self.name = name
# 实现 run() 方法,定义进程的业务逻辑
def run(self):
print(f"进程 {self.name} 运行")
time.sleep(2)
print(f"进程 {self.name} 结束")
if __name__ == "__main__":
p1 = MyProcess("进程1")
p2 = MyProcess("进程2")
p1.start()
p2.start()
p1.join()
p2.join()运行示例,输出如下:
进程 进程2 运行
进程 进程2 结束
进程 进程1 运行
进程 进程1 结束注意了,如果在 Windows 中,没有将创建进程的代码包在 if __name__ == "__main__": 中,会出现如下错误:
Traceback (most recent call last):
File "<string>", line 1, in <module>
from multiprocessing.spawn import spawn_main; spawn_main(parent_pid=17280, pipe_handle=312)
~~~~~~~~~~^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "D:\Program Files\Python313\Lib\multiprocessing\spawn.py", line 122, in spawn_main
exitcode = _main(fd, parent_sentinel)
.......
RuntimeError:
An attempt has been made to start a new process before the
current process has finished its bootstrapping phase.
This probably means that you are not using fork to start your
child processes and you have forgotten to use the proper idiom
in the main module:
if __name__ == '__main__':
freeze_support()
...
The "freeze_support()" line can be omitted if the program
is not going to be frozen to produce an executable.
To fix this issue, refer to the "Safe importing of main module"
section in https://docs.python.org/3/library/multiprocessing.html上面的 RuntimeError 错误是 Windows 系统使用 multiprocessing 经典报错:
Windows 创建进程默认用 spawn 方式,子进程会重新导入你的整个脚本文件。如果创建进程 p.start() 代码直接写在全局作用域(没有包裹 if __name__ == "__main__":)。子进程导入脚本时,又会再次执行 p1.start(),无限递归创建进程,Python 直接抛出异常拦截。
但是,Linux/WSL 默认 fork 不会出现该问题,但 Windows、打包 exe 必须遵守该规范。
一次性创建多个进程,重复使用,避免频繁创建销毁开销,可以使用 multiprocessing.Pool。例如:
from multiprocessing import Pool, current_process
import time
import os
def task(num):
# 获取当前执行任务的进程信息
proc = current_process()
worker_name = proc.name # 进程池worker名称:ForkPoolWorker-x
pid = os.getpid() # 操作系统PID
print(f"任务 {num} 执行 | 进程名: {worker_name}, PID: {pid}")
time.sleep(1)
return f"任务 {num} 完成 | 进程PID:{pid}"
if __name__ == "__main__":
# 创建进程池,最多同时运行 3 个进程
pool = Pool(3)
# 异步执行 10 个任务
results = []
for i in range(10):
res = pool.apply_async(task, args=(i,))
results.append(res)
# 关闭进程池,不再接受新任务
pool.close()
# 等待所有任务完成
pool.join()
# 获取结果
for r in results:
print(r.get())运行示例,输出如下:
任务 1 执行 | 进程名: SpawnPoolWorker-2, PID: 12576
任务 4 执行 | 进程名: SpawnPoolWorker-2, PID: 12576
任务 7 执行 | 进程名: SpawnPoolWorker-2, PID: 12576
任务 2 执行 | 进程名: SpawnPoolWorker-3, PID: 21188
任务 5 执行 | 进程名: SpawnPoolWorker-3, PID: 21188
任务 8 执行 | 进程名: SpawnPoolWorker-3, PID: 21188
任务 0 执行 | 进程名: SpawnPoolWorker-1, PID: 25852
任务 3 执行 | 进程名: SpawnPoolWorker-1, PID: 25852
任务 6 执行 | 进程名: SpawnPoolWorker-1, PID: 25852
任务 9 执行 | 进程名: SpawnPoolWorker-1, PID: 25852
任务 0 完成 | 进程PID:25852
任务 1 完成 | 进程PID:12576
任务 2 完成 | 进程PID:21188
任务 3 完成 | 进程PID:25852
任务 4 完成 | 进程PID:12576
任务 5 完成 | 进程PID:21188
任务 6 完成 | 进程PID:25852
任务 7 完成 | 进程PID:12576
任务 8 完成 | 进程PID:21188
任务 9 完成 | 进程PID:25852使用进程池能够实现高效任务执行,可灵活控制并发数量,整体任务管理流程简单便捷。
进程内存独立,不能直接共享数据,如果要互相传递数据,必须通过通信机制。
Python 中进程通信的常用方式:
(1)Queue 队列(进程安全)
(2)Pipe 管道
(3)共享内存
下面将演示如何通过进程队列 Queue 进行通信,例如:
from multiprocessing import Process, Queue
import time
# 写入数据到队列的进程
def write(q):
for i in ["A", "B", "C"]:
q.put(i)
print(f"放入:{i}")
time.sleep(0.5)
# 从队列消费数据的进程
def read(q):
while True:
if not q.empty():
data = q.get()
print(f"取出:{data}")
time.sleep(0.5)
else:
break
if __name__ == "__main__":
# 创建队列
q = Queue()
# 创建进程,将队列作为参数
p1 = Process(target=write, args=(q,))
p2 = Process(target=read, args=(q,))
p1.start()
p1.join()
p2.start()
p2.join()运行程序,输出如下:
放入:A
放入:B
放入:C
取出:A
取出:B
取出:C该示例将完整模拟一套售票程序的运行逻辑与执行流程。目的是了解多进程的运用和协调。
多个进程同时执行卖票操作时,会出现并发竞争问题,因此必须通过加锁机制约束并发访问,以此保障票务相关数据的一致性与数据安全。
代码如下:
from multiprocessing import Process, Lock, Value
import time
# 售票的进程
def sell(name, ticket, lock):
while True:
with lock:
if ticket.value > 0:
ticket.value -= 1
print(f"{name} 卖出1张,剩余:{ticket.value}")
time.sleep(0.2)
else:
break
if __name__ == "__main__":
ticket = Value("i", 10) # 共10张票
lock = Lock() # 全局锁
# 启动两个进程
p1 = Process(target=sell, args=("窗口1", ticket, lock))
p2 = Process(target=sell, args=("窗口2", ticket, lock))
# 启动进程
p1.start()
p2.start()
# 等待进程结束
p1.join()
p2.join()
print("票已售罄")运行示例,输出如下:
窗口1 卖出1张,剩余:9
窗口1 卖出1张,剩余:7
窗口1 卖出1张,剩余:5
窗口1 卖出1张,剩余:3
窗口1 卖出1张,剩余:1
窗口2 卖出1张,剩余:8
窗口2 卖出1张,剩余:6
窗口2 卖出1张,剩余:4
窗口2 卖出1张,剩余:2
窗口2 卖出1张,剩余:0
票已售罄注意了:
Value("i", 10) 语句将在共享内存中创建变量,所有子进程读写同一份数据,而不是各自拷贝副本。访问必须用 .value 语法。例如,ticket.value。注意,不要使用普通全局变量创建 ticket,创建子进程时,子进程会拷贝一份独立的 ticket=10。窗口 1 操作自己的 ticket,窗口 2 操作另一份 ticket,互不影响,最终各自卖 10 张。
即使使用 Value("i", 10) 也不要放在全局变量上面,Windows 系统 multiprocessing 创建子进程使用 spawn 方式,子进程会重新导入整个模块,重新执行顶层代码,最终每个子进程内部又新建了一份独立的 ticket。
当前示例演示在进程内部创建并启动线程:进程依托计算机多核硬件资源实现算力利用,创建出的线程则负责承担各类并发 I/O 相关任务处理工作。代码如下:
from multiprocessing import Process
from threading import Thread
import time
# 线程任务
def thread_task(name):
print(f"线程 {name} 运行")
time.sleep(1)
# 进程任务
def process_task(name):
print(f"进程 {name} 启动")
# 进程内启动 2 个线程
t1 = Thread(target=thread_task, args=(f"{name}-T1",))
t2 = Thread(target=thread_task, args=(f"{name}-T2",))
t1.start()
t2.start()
t1.join()
t2.join()
print(f"进程 {name} 结束")
if __name__ == "__main__":
# 启动两个进程
p1 = Process(target=process_task, args=("P1",))
p2 = Process(target=process_task, args=("P2",))
p1.start()
p2.start()
p1.join()
p2.join()
print("全部完成")运行实例,输出如下:
进程 P1 启动
线程 P1-T1 运行
线程 P1-T2 运行
进程 P1 结束
进程 P2 启动
线程 P2-T1 运行
线程 P2-T2 运行
进程 P2 结束
全部完成更多并发信息参考 https://docs.python.org/zh-cn/3.11/library/concurrency.html 文档。