Когда вы передаете большие данные между процессами через multiprocessing.Queue, каждое сообщение копируется, сериализуется через pickle и блокирует GIL. Для видео, аудио или numpy-массивов это становится узким горлом. Многие разработчики ошибочно уходят в asyncio в одном процессе, теряя возможность использовать многоядерность.
Идея zero-copy через shared memory
Используйте multiprocessing.shared_memory для выделения пула буферов. Данные кладутся без копирования, а очередь строится как кольцевой буфер с атомарными head и tail через ctypes.Value. Это SPMC (single producer, multiple consumers) без мьютексов.
Пример реализации ring buffer
from multiprocessing import shared_memory, Value
import numpy as np
class RingBuffer:
def __init__(self, size=4, shm_name='my_shm'):
self.size = size
self.buf = shared_memory.SharedMemory(
name=shm_name, create=True, size=4096)
self.head = Value('i', 0)
self.tail = Value('i', 0)
def push(self, data: np.ndarray):
offset = self.head.value * 1024
self.buf.buf[offset:offset + len(data.tobytes())] = data.tobytes()
self.head.value = (self.head.value + 1) % self.size
def pop(self) -> np.ndarray:
while self.tail.value == self.head.value:
pass # spinlock (production: use futex)
offset = self.tail.value * 1024
raw = self.buf.buf[offset:offset + 1024]
self.tail.value = (self.tail.value + 1) % self.size
return np.frombuffer(raw, dtype=np.uint8)
Типичная ошибка и trade-off
Основная ошибка - игнорирование фиксированного размера сообщений и отсутствие защиты от перезаписи. Если producer запишет больше данных, чем буфер, данные молча перетрутся. Также spinlock в pop() нагружает CPU - для production используйте pthread_cond_wait через ctypes.
Когда это оправдано
Для тяжелых данных: видеофреймы, аудио, большие матрицы в ETL или ML-инференсе. Выигрыш по сравнению с multiprocessing.Queue + pickle может достигать 3-10x на больших блоках. Для мелких сообщений overhead shared memory не оправдан.
Вывод: Lock-free очередь на shared_memory с zero-copy дает радикальный прирост производительности для больших данных между процессами, но требует аккуратного управления размером буфера и синхронизации.