Python3 多进程编程

🎉摘要:本文深入讲解Python多进程编程,包括使用multiprocessing模块创建进程(函数式与继承类)、进程池Pool管理并发任务、进程间通信(Queue队列),并通过模拟售票等实战案例展示如何规避GIL限制,实现CPU密集型任务的并行加速。

前面介绍了多线程编程,由于 Python 中,CPython 存在 GIL(全局解释器锁),同一个进程内,同一时刻仅有一个线程执行 Python 字节码。CPU 密集型场景,多线程无法利用多核并行,提速效果极差,甚至比单线程更慢。

为了解决 CPU 密集任务并发问题,推出了多进程编程。多进程可以创建独立的子进程,每个进程有自己独立的内存空间,完全并行,不受 GIL(全局解释器锁)限制。

os.fork() 函数

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 模块

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 类

创建一个自己的类,继承 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 文档。

说说我的看法
全部评论(
没有评论
关于
本网站专注于 Java、数据库(MySQL、Oracle)、Linux、软件架构及大数据等多领域技术知识分享。涵盖丰富的原创与精选技术文章,助力技术传播与交流。无论是技术新手渴望入门,还是资深开发者寻求进阶,这里都能为您提供深度见解与实用经验,让复杂编码变得轻松易懂,携手共赴技术提升新高度。如有侵权,请来信告知:hxstrive@outlook.com
其他应用
公众号