Python 多进程编程:multiprocessing 模块的进程间通信技巧
·
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 | 复杂数据结构共享 | 低 | 低 |
最佳实践建议
- 避免共享状态:优先使用消息传递(队列/管道)而非共享内存
- 处理僵尸进程:始终调用
join()或设置daemon=True - 异常处理:使用
try/except捕获Queue.Empty等异常 - 性能优化:
- 批量处理数据减少通信次数
- 使用
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]
更多推荐


所有评论(0)