Python 多进程编程:multiprocessing 模块的进程间通信技巧

在多进程编程中,进程间通信(IPC)是关键挑战。multiprocessing 模块提供以下核心通信机制:


1. 队列(Queue)

最常用的线程安全通信方式,支持生产者-消费者模式:

from multiprocessing import Process, Queue

def worker(q):
    q.put("子进程数据")  # 写入队列

if __name__ == "__main__":
    q = Queue()
    p = Process(target=worker, args=(q,))
    p.start()
    print(q.get())  # 输出: "子进程数据"
    p.join()

特性:

  • 自动处理锁机制
  • 支持 put()/get() 阻塞操作
  • 可通过 maxsize 限制队列长度

2. 管道(Pipe)

双向通信通道,适合一对一进程通信:

from multiprocessing import Process, Pipe

def worker(conn):
    conn.send("消息A")  # 发送数据
    print(conn.recv())  # 接收数据

if __name__ == "__main__":
    parent_conn, child_conn = Pipe()
    p = Process(target=worker, args=(child_conn,))
    p.start()
    print(parent_conn.recv())  # 输出: "消息A"
    parent_conn.send("消息B")  # 发送响应
    p.join()

特性:

  • 返回两个连接对象(双向)
  • send()/recv() 方法实现数据传输
  • 需注意死锁风险(双方同时调用 recv()

3. 共享内存(Value/Array)

直接共享内存空间,适用于高性能场景:

from multiprocessing import Process, Value, Array

def worker(n, arr):
    n.value *= 2  # 修改共享值
    for i in range(len(arr)):
        arr[i] **= 2  # 修改共享数组

if __name__ == "__main__":
    num = Value("i", 5)  # 整型共享内存
    arr = Array("d", [1.0, 2.0, 3.0])  # 双精度数组

    p = Process(target=worker, args=(num, arr))
    p.start()
    p.join()

    print(num.value)  # 输出: 10
    print(arr[:])     # 输出: [1.0, 4.0, 9.0]

特性:

  • Value(typecode, value):共享单个值
  • Array(typecode, sequence):共享数组
  • 需用锁(Lock)保证原子性:
    num = Value("i", 0)
    lock = Lock()
    with lock:
        num.value += 1
    


4. 管理器(Manager)

托管共享对象,支持复杂数据结构:

from multiprocessing import Process, Manager

def worker(d, l):
    d["key"] = "value"  # 修改共享字典
    l.append(10)        # 修改共享列表

if __name__ == "__main__":
    with Manager() as manager:
        d = manager.dict()  # 共享字典
        l = manager.list()  # 共享列表

        p = Process(target=worker, args=(d, l))
        p.start()
        p.join()

        print(d)  # 输出: {'key': 'value'}
        print(l)  # 输出: [10]

适用场景:

  • 共享字典、列表等复杂结构
  • 跨进程同步数据库连接池
  • 网络套接字共享

通信机制选择指南
机制适用场景性能复杂度
Queue生产者-消费者模型
Pipe一对一实时通信
共享内存大数据量/高频操作极高
Manager复杂数据结构共享

最佳实践建议
  1. 避免共享状态:优先使用消息传递(队列/管道)而非共享内存
  2. 处理僵尸进程:始终调用 join() 或设置 daemon=True
  3. 异常处理:使用 try/except 捕获 Queue.Empty 等异常
  4. 性能优化
    • 批量处理数据减少通信次数
    • 使用 Queue 替代 Manager 提升速度
    • 共享内存时用 RawArray 避免锁开销

示例:高效数据并行处理

from multiprocessing import Pool

def process_data(data):
    return data * 2  # 模拟计算

if __name__ == "__main__":
    with Pool(4) as p:  # 4进程池
        results = p.map(process_data, range(10))  # 自动分配任务
    print(results)  # 输出: [0, 2, 4, ..., 18]

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐