Python脚本多进程如何共享数据

wen 实用脚本 9

本文目录导读:

Python脚本多进程如何共享数据

  1. multiprocessing.Value / Array(适用于简单数据类型)
  2. multiprocessing.Manager(更灵活,支持复杂数据结构)
  3. multiprocessing.Queue(线程安全的队列)
  4. multiprocessing.Pipe(双向通信)
  5. 共享内存(推荐新代码使用)
  6. 选择建议
  7. 注意事项

在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 大型数组/数据块 最高

注意事项

  1. 锁机制:多个进程同时修改共享数据时,需要使用锁:
    from multiprocessing import Lock

lock = Lock() with lock: shared_value.value += 1


2. **避免使用全局变量**:进程间不共享全局变量
3. **序列化开销**:Manager和Queue需要序列化数据,大数据量时性能下降
4. **资源管理**:记得正确关闭和释放共享资源

抱歉,评论功能暂时关闭!