目录

Python 多进程编程

学习目标

  • 理解进程与线程的区别
  • 掌握 multiprocessing 模块创建进程
  • 学会进程间通信(Queue、Pipe、共享内存)
  • 掌握进程池的使用

1. 进程基础

1.1 创建进程

import multiprocessing
import os
import time

def worker(name):
    """工作函数"""
    print(f"进程 {name} (PID: {os.getpid()}) 开始")
    time.sleep(2)
    print(f"进程 {name} (PID: {os.getpid()}) 结束")

if __name__ == '__main__':
    print(f"主进程 PID: {os.getpid()}")
    
    # 创建进程
    p1 = multiprocessing.Process(target=worker, args=("A",))
    p2 = multiprocessing.Process(target=worker, args=("B",))
    
    p1.start()
    p2.start()
    
    p1.join()
    p2.join()
    
    print("所有进程完成")

1.2 继承 Process 类

import multiprocessing
import time

class MyProcess(multiprocessing.Process):
    def __init__(self, name, delay):
        super().__init__()
        self.name = name
        self.delay = delay
    
    def run(self):
        print(f"进程 {self.name} 开始")
        time.sleep(self.delay)
        print(f"进程 {self.name} 结束")

if __name__ == '__main__':
    processes = [
        MyProcess("工作1", 1),
        MyProcess("工作2", 2),
    ]
    
    for p in processes:
        p.start()
    for p in processes:
        p.join()

2. 进程间通信

2.1 Queue 队列

import multiprocessing
import time

def producer(queue):
    for i in range(5):
        item = f"产品-{i}"
        queue.put(item)
        print(f"生产: {item}")
        time.sleep(0.5)

def consumer(queue):
    while True:
        try:
            item = queue.get(timeout=3)
            print(f"消费: {item}")
        except:
            print("超时退出")
            break

if __name__ == '__main__':
    queue = multiprocessing.Queue(maxsize=10)
    
    p = multiprocessing.Process(target=producer, args=(queue,))
    c = multiprocessing.Process(target=consumer, args=(queue,))
    
    p.start()
    c.start()
    p.join()
    c.join()

2.2 Pipe 管道

import multiprocessing

def sender(conn):
    messages = ["hello", "world", "done"]
    for msg in messages:
        conn.send(msg)
        print(f"发送: {msg}")
    conn.close()

def receiver(conn):
    while True:
        try:
            msg = conn.recv()
            print(f"接收: {msg}")
            if msg == "done":
                break
        except EOFError:
            break

if __name__ == '__main__':
    # 创建管道
    parent_conn, child_conn = multiprocessing.Pipe()
    
    p = multiprocessing.Process(target=sender, args=(child_conn,))
    p.start()
    
    receiver(parent_conn)
    p.join()

2.3 共享内存

import multiprocessing

def increment(counter, lock):
    for _ in range(10000):
        with lock:
            counter.value += 1

if __name__ == '__main__':
    # 共享变量
    counter = multiprocessing.Value('i', 0)  # 'i' = 整数
    lock = multiprocessing.Lock()
    
    processes = [
        multiprocessing.Process(target=increment, args=(counter, lock))
        for _ in range(5)
    ]
    
    for p in processes:
        p.start()
    for p in processes:
        p.join()
    
    print(f"最终计数: {counter.value}")  # 50000

2.4 Manager 共享对象

import multiprocessing

def worker(d, l, n):
    d[n] = n * n
    l.append(n)
    print(f"进程 {n}: dict={dict(d)}, list={list(l)}")

if __name__ == '__main__':
    with multiprocessing.Manager() as manager:
        shared_dict = manager.dict()
        shared_list = manager.list()
        
        processes = [
            multiprocessing.Process(target=worker, args=(shared_dict, shared_list, i))
            for i in range(5)
        ]
        
        for p in processes:
            p.start()
        for p in processes:
            p.join()
        
        print(f"最终 dict: {dict(shared_dict)}")
        print(f"最终 list: {list(shared_list)}")

3. 进程池

3.1 Pool 基本使用

import multiprocessing
import time

def square(n):
    time.sleep(0.5)
    return n * n

if __name__ == '__main__':
    numbers = [1, 2, 3, 4, 5]
    
    # 创建进程池
    with multiprocessing.Pool(processes=3) as pool:
        # map - 按顺序返回结果
        results = pool.map(square, numbers)
        print(f"map 结果: {results}")
        
        # apply_async - 异步提交单个任务
        result = pool.apply_async(square, (10,))
        print(f"async 结果: {result.get()}")
        
        # imap - 迭代返回结果
        for result in pool.imap(square, numbers):
            print(f"imap: {result}")
        
        # imap_unordered - 按完成顺序返回
        for result in pool.imap_unordered(square, numbers):
            print(f"unordered: {result}")

3.2 进程池对比

import multiprocessing
import time

def cpu_task(n):
    """CPU 密集型任务"""
    count = 0
    for i in range(n):
        count += i * i
    return count

if __name__ == '__main__':
    tasks = [1000000] * 4
    
    # 串行执行
    start = time.time()
    serial_results = [cpu_task(t) for t in tasks]
    print(f"串行: {time.time() - start:.2f}s")
    
    # 多进程
    start = time.time()
    with multiprocessing.Pool() as pool:
        parallel_results = pool.map(cpu_task, tasks)
    print(f"多进程: {time.time() - start:.2f}s")

4. 进程同步

4.1 Lock

import multiprocessing

def task(lock, shared_list, item):
    with lock:
        shared_list.append(item)
        print(f"添加 {item},列表: {list(shared_list)}")

if __name__ == '__main__':
    with multiprocessing.Manager() as manager:
        lock = multiprocessing.Lock()
        shared_list = manager.list()
        
        processes = [
            multiprocessing.Process(target=task, args=(lock, shared_list, i))
            for i in range(5)
        ]
        
        for p in processes:
            p.start()
        for p in processes:
            p.join()

4.2 Semaphore

import multiprocessing
import time

def limited_task(semaphore, name):
    with semaphore:
        print(f"{name} 开始执行")
        time.sleep(2)
        print(f"{name} 执行完成")

if __name__ == '__main__':
    # 最多允许 2 个进程同时执行
    semaphore = multiprocessing.Semaphore(2)
    
    processes = [
        multiprocessing.Process(target=limited_task, args=(semaphore, f"任务{i}"))
        for i in range(5)
    ]
    
    for p in processes:
        p.start()
    for p in processes:
        p.join()

5. 进程 vs 线程

特性 进程 线程
内存空间 独立 共享
创建开销
通信方式 IPC(复杂) 直接共享(简单)
GIL 影响 无影响 受限制
适用场景 CPU 密集型 IO 密集型
崩溃影响 不影响其他进程 可能导致整个程序崩溃
import multiprocessing
import threading
import time

def cpu_bound(n):
    """CPU 密集型"""
    count = 0
    for i in range(n):
        count += i * i
    return count

def io_bound(delay):
    """IO 密集型"""
    time.sleep(delay)
    return "完成"

if __name__ == '__main__':
    # CPU 密集型对比
    print("=== CPU 密集型 ===")
    
    start = time.time()
    with multiprocessing.Pool(4) as pool:
        pool.map(cpu_bound, [5000000] * 4)
    print(f"多进程: {time.time() - start:.2f}s")
    
    start = time.time()
    threads = [threading.Thread(target=cpu_bound, args=(5000000,)) for _ in range(4)]
    for t in threads: t.start()
    for t in threads: t.join()
    print(f"多线程: {time.time() - start:.2f}s")
    
    # IO 密集型对比
    print("\n=== IO 密集型 ===")
    
    start = time.time()
    with multiprocessing.Pool(4) as pool:
        pool.map(io_bound, [1] * 10)
    print(f"多进程: {time.time() - start:.2f}s")
    
    start = time.time()
    threads = [threading.Thread(target=io_bound, args=(1,)) for _ in range(10)]
    for t in threads: t.start()
    for t in threads: t.join()
    print(f"多线程: {time.time() - start:.2f}s")

本节小结

  • 创建进程Process(target=...) 或继承 Process
  • 进程通信QueuePipeValueArrayManager
  • 进程池Pool 简化多进程管理,适合批量任务
  • 同步LockSemaphore 保护共享资源
  • 选择:CPU 密集型用多进程,IO 密集型用多线程

练习

  1. 使用多进程并行计算大列表中每个元素的平方
  2. 实现一个多进程文件处理器,每个进程处理一个文件
  3. 使用 Manager 实现一个跨进程的字典缓存
  4. 比较多进程和多线程处理 CPU 密集型任务的性能