При передаче большого объекта через multiprocessing.Queue почему память временно возрастает у отправляющего процесса?
multiprocessing.Queue обычно не передаёт объект напрямую: он сериализуется через pickle, а сериализованные данные буферизуются для отправки в другой процесс. Поэтому некоторое время у отправителя могут одновременно находиться исходный объект и его сериализованное представление, что увеличивает пиковое потребление памяти.
Получатель затем создаёт отдельный объект при десериализации. Передача через очередь не является механизмом совместного доступа к одной области памяти.
Модуль multiprocessing изолирует процессы, чтобы они имели независимые адресные пространства и могли использовать несколько ядер без общего состояния интерпретатора. Такая изоляция повышает надёжность границ между задачами, но требует явного обмена данными.
multiprocessing.Queue предоставляет удобный обмен Python-объектами, поэтому использует сериализацию. Это проще, чем вручную управлять общей памятью, но требует дополнительных CPU-операций и памяти.
Если в очередь помещаются большие изображения, таблицы или бинарные буферы, отправитель может удерживать исходный объект, сериализованный буфер и внутренние буферы очереди. Получатель дополнительно выделяет память под восстановленный объект.
Пиковое потребление может быть значительно выше размера исходных данных. Неправильная оценка особенно опасна при пакетной отправке: несколько ещё не обработанных сообщений способны накопиться в очереди.
Даже при использовании fork помещение объекта в очередь не становится автоматически безкопийным. Наследованная память может использоваться через copy-on-write, но обмен новым объектом через очередь всё равно обычно проходит через сериализацию.
При вызове put объект передаётся внутреннему механизму очереди. В реализации multiprocessing.Queue отдельный поток-обработчик сериализует объект и отправляет байты через канал между процессами. До завершения передачи исходный объект может оставаться живым, а сериализованное представление — занимать дополнительную память.
На стороне получателя байты распаковываются, и создаётся новый объект. Для обычных составных объектов это означает как минимум отдельные структуры контейнеров и часто отдельные копии вложенных данных.
Минимальная иллюстрация механизма:
bytearray остаётся у отправителя после put, пока на него существуют ссылки, а очередь передаёт его сериализованное содержимое. Кроме того, put может завершиться до фактической записи всех данных благодаря внутреннему буферу, поэтому момент возврата из put не означает завершение передачи.
Для снижения пикового потребления применяют следующие подходы:
mmap, когда данные естественно представлены файлом;Общая память уменьшает копирование, но усложняет синхронизацию, управление временем жизни и согласованность данных. Потоки устраняют сериализацию внутри одного процесса, однако не дают изоляции процессов и не всегда ускоряют CPU-bound Python-код.
Сервис обрабатывает изображения пакетами. Сначала один процесс читал изображения и помещал целые массивы в multiprocessing.Queue. При нескольких крупных заданиях память резко возрастала: отправитель удерживал исходные массивы, очередь — сериализованные сообщения, а рабочие процессы создавали свои экземпляры после распаковки.
Рассматривались три варианта. Потоки были проще, но не обеспечивали нужного ускорения CPU-bound обработки. Ограничение размера очереди уменьшало пик, но снижало пропускную способность. Передача через общую память требовала явного протокола владения буферами, зато позволяла передавать дескриптор области памяти вместо самого массива.
Выбрали общую память с небольшой очередью метаданных: через очередь передавались имя области, размер и параметры изображения, а не сами пиксели. Это снизило лишние копирования и сделало пик памяти предсказуемее; взамен добавили контроль освобождения областей и защиту от преждевременного закрытия буфера.
put завершение передачи объекта?Нет. Для multiprocessing.Queue put обычно лишь помещает объект во внутреннюю структуру очереди, после чего фоновый механизм выполняет сериализацию и отправку. Поэтому измерение памяти сразу после put может не показывать уже завершённое состояние, а временный пик может возникнуть позже.
fork копирование больших данных при передаче их рабочему процессу?fork позволяет дочернему процессу первоначально видеть страницы памяти родителя без немедленного копирования благодаря copy-on-write. Но это относится к наследованию адресного пространства при создании процесса. Если объект затем передаётся через Queue, он обычно сериализуется отдельно; если дочерний процесс изменяет унаследованные страницы, возникают дополнительные физические копии.
Нет. Она может убрать сериализованные копии, но сама область общей памяти занимает место, а рабочий код может создать дополнительные временные представления или преобразования. Кроме того, при ошибках протокола области могут не освобождаться, а одновременная обработка нескольких буферов всё равно увеличит пик. Общую память нужно сочетать с ограничением числа буферов и ясным управлением их временем жизни.