本文目录导读:

- multiprocessing.Value / Array(适用于简单数据类型)
- multiprocessing.Manager(更灵活,支持复杂数据结构)
- multiprocessing.Queue(线程安全的队列)
- multiprocessing.Pipe(双向通信)
- 共享内存(推荐新代码使用)
- 选择建议
- 注意事项
在Python多进程中共享数据,主要有以下几种常用方法:
multiprocessing.Value / Array(适用于简单数据类型)
from multiprocessing import Process, Value, Array
def worker(val, arr):
val.value += 1
for i in range(len(arr)):
arr[i] *= 2
if __name__ == '__main__':
shared_value = Value('i', 0) # 'i' 表示整数类型
shared_array = Array('d', [1.0, 2.0, 3.0]) # 'd' 表示双精度浮点数
p = Process(target=worker, args=(shared_value, shared_array))
p.start()
p.join()
print(shared_value.value)
print(shared_array[:])
multiprocessing.Manager(更灵活,支持复杂数据结构)
from multiprocessing import Process, Manager
def worker(shared_dict, shared_list):
shared_dict['key'] = 'value'
shared_list.append('new_item')
if __name__ == '__main__':
with Manager() as manager:
shared_dict = manager.dict()
shared_list = manager.list()
p = Process(target=worker, args=(shared_dict, shared_list))
p.start()
p.join()
print(shared_dict)
print(shared_list)
multiprocessing.Queue(线程安全的队列)
from multiprocessing import Process, Queue
def producer(q):
for i in range(5):
q.put(f'message {i}')
def consumer(q):
while True:
msg = q.get()
if msg is None: # 结束信号
break
print(f'Received: {msg}')
if __name__ == '__main__':
q = Queue()
p1 = Process(target=producer, args=(q,))
p2 = Process(target=consumer, args=(q,))
p1.start()
p2.start()
p1.join()
q.put(None) # 发送结束信号
p2.join()
multiprocessing.Pipe(双向通信)
from multiprocessing import Process, Pipe
def child_process(conn):
data = conn.recv()
print(f"Child received: {data}")
conn.send("Hello from child!")
if __name__ == '__main__':
parent_conn, child_conn = Pipe()
p = Process(target=child_process, args=(child_conn,))
p.start()
parent_conn.send("Hello from parent!")
response = parent_conn.recv()
print(f"Parent received: {response}")
p.join()
共享内存(推荐新代码使用)
Python 3.8+ 提供更高效的共享内存:
from multiprocessing import shared_memory, Process
import numpy as np
def worker(shm_name, shape, dtype):
existing_shm = shared_memory.SharedMemory(name=shm_name)
arr = np.ndarray(shape, dtype=dtype, buffer=existing_shm.buf)
arr[:] = arr * 2
if __name__ == '__main__':
data = np.array([1, 2, 3, 4])
shm = shared_memory.SharedMemory(create=True, size=data.nbytes)
shared_arr = np.ndarray(data.shape, dtype=data.dtype, buffer=shm.buf)
shared_arr[:] = data[:]
p = Process(target=worker, args=(shm.name, data.shape, data.dtype))
p.start()
p.join()
print(shared_arr)
shm.close()
shm.unlink()
选择建议
| 方法 | 适用场景 | 性能 | 复杂度 |
|---|---|---|---|
| Value/Array | 简单数据类型共享 | 高 | 低 |
| Manager | 复杂数据结构 | 中 | 中 |
| Queue | 生产者-消费者模式 | 中 | 中 |
| Pipe | 双向通信 | 高 | 中 |
| Shared Memory | 大型数组/数据块 | 最高 | 高 |
注意事项
- 锁机制:多个进程同时修改共享数据时,需要使用锁:
from multiprocessing import Lock
lock = Lock() with lock: shared_value.value += 1
2. **避免使用全局变量**:进程间不共享全局变量
3. **序列化开销**:Manager和Queue需要序列化数据,大数据量时性能下降
4. **资源管理**:记得正确关闭和释放共享资源