Memory not deallocated when using pybind11 with multiprocessing and shared memory

Viewed 236

I am using pybind11 to create a class called Object (defined in obj.h and obj.cpp and bound from pybind in bindings.cpp), which can have a large vector added to it (3,000,000 string objects), and can return an edited vector to python in a series of iterations that are parallelised.

I place Object in shared memory (shared_memory.py) via a numpy array, and then call return_vec method from within a multiprocessing loop (__main__.py, main()).

I have noticed using psutil that the memory is not deallocated between each iteration of the loop, despite the python function (__main__.py, addition_wrapper()) itself not returning anything. Setting maxtasksperchild=1 deallocates the memory between each iteration. The object's copy constructor is not called, so the memory allocation is not coming from repeated copying.

Would someone be able to explain what is occurring, and potentially provide a solution if this is a memory leak, please?

Thanks in advance for any help you can provide.

Results (all results are in % of physical memory allocated to the python subprocess):

Single process and maxtasksperchild=None

Pre-construction: 0.18
Object default constructor
Post-construction: 0.89
p0 pre-addition: 0.02
p0 post-addition: 2.15
p1 pre-addition: 2.01
p1 post-addition: 2.15
p2 pre-addition: 2.15
p2 post-addition: 2.15
p3 pre-addition: 2.15
p3 post-addition: 2.15
p4 pre-addition: 2.15
p4 post-addition: 2.15
Post-iteration: 0.89

Single process and maxtasksperchild=1

Pre-construction: 0.18
Object default constructor
Post-construction: 0.89
p0 pre-addition: 0.02
p0 post-addition: 2.15
p1 pre-addition: 0.03
p1 post-addition: 2.15
p2 pre-addition: 0.02
p2 post-addition: 2.15
p3 pre-addition: 0.02
p3 post-addition: 2.15
p4 pre-addition: 0.02
p4 post-addition: 2.15
Post-iteration: 0.89

Scripts:

obj.h

#ifndef OBJ_H
#define OBJ_H

// STL headers
#include <iostream>
#include <fstream>
#include <vector>
#include <string>

// pybind11 headers
#include <pybind11/pybind11.h>
#include <pybind11/stl.h>

// global variable declaration
namespace py = pybind11;

class Object {
    public:
            // constructors
    // default constructor
    Object()
    {
        std::cout << "Object default constructor" << std::endl;
    };

    // copy
    Object(const Object & rhs)
    {
          std::cout << "Object copy constructor" << std::endl;
          _private_vec = rhs._private_vec;
    };
    // copy operator
    Object & operator = (const Object & rhs)
    {
          std::cout << "Object copy operator" << std::endl;
          if(this != &rhs)
          {
            _private_vec = rhs._private_vec;
          }
          return *this;
    };

    // move
    Object(const Object && rhs) noexcept
    {
          std::cout << "Object move constructor" << std::endl;
          _private_vec = std::move(rhs._private_vec);
    };
    // move operator
    Object & operator = (const Object && rhs) noexcept
    {
          std::cout << "Object move operator" << std::endl;
          if(this != &rhs)
          {
            _private_vec = std::move(rhs._private_vec);
          }
          return *this;
    };

    // destructor
    ~Object()
    {
        std::cout << "Object destructor" << std::endl;
    };
    
    // Object methods

    // add vector to object
    void add(const std::vector<std::string>& string_vec);

    // return vector after adding string
    std::vector<std::string> return_vec(const std::string& addition);

    private:
    // private vector object
    std::vector<std::string> _private_vec;
};

#endif //OBJ_H

obj.cpp

#include "obj.h"

void Object::add(const std::vector<std::string>& string_vec)
{
    _private_vec = string_vec;
}

std::vector<std::string> Object::return_vec(const std::string& addition)
{
    std::vector<std::string> updated_vec(_private_vec.size());

    for (int i = 0; i < _private_vec.size(); i++)
    {
        const auto new_str = _private_vec[i] + addition;
        updated_vec[i] = new_str;
    }

    return updated_vec;
}

bindings.cpp

// ggCaller header
#include "bindings.h"

PYBIND11_MODULE(test_cpp, m)
{
    m.doc() = "Simple vector updating";

    py::class_<Object, std::shared_ptr<Object>>(m, "Graph")
            .def("add", &Object::add)
            .def("return_vec", &Object::return_vec);

    m.def("create_obj", []() { return std::make_shared<Object>(); });
}

__main__.py

import argparse
from test_pybind.shared_memory import *
import test_cpp
import psutil
import numpy as np
import tqdm
from functools import partial

def get_options():
    description = 'Runs simple test of cpp bound functions'
    parser = argparse.ArgumentParser(description=description,
                                     prog='test_pybind')

    IO = parser.add_argument_group('Input/Output options')
    IO.add_argument('--str',
                    type=str,
                    default=None,
                    help='String to add to function ')
    IO.add_argument('--num',
                    type=int,
                    default=None,
                    help='Number of copies of string')
    IO.add_argument('--iter',
                    type=int,
                    default=100,
                    help='Number of iterations')
    IO.add_argument('--maxtasks',
                    type=int,
                    default=None,
                    help='Number of tasks per child process')
    IO.add_argument('--threads',
                    type=int,
                    default=1,
                    help='Number threads to run')

    return parser.parse_args()


def addition_wrapper(addition, array_shd_tup):
    iteration, addition_str = addition

    p = psutil.Process()
    print("p" + str(iteration) + " pre-addition: " + str(p.memory_percent()))

    # unpack array
    existing_shm = shared_memory.SharedMemory(name=array_shd_tup.name)
    shd_arr = np.ndarray(array_shd_tup.shape, dtype=array_shd_tup.dtype, buffer=existing_shm.buf)

    print("p" + str(iteration) + " Object ID: " + str(id(shd_arr[0])))

    # return newly appended vector (shd_arr[0] is cpp_obj)
    new_vec = shd_arr[0].return_vec(addition_str)

    print("p" + str(iteration) + " post-addition: " + str(p.memory_percent()))

    return


def main():
    #psutils memory tracker
    p0 = psutil.Process()

    print("Pre-construction: " + str(p0.memory_percent()))

    options = get_options()

    # create new object
    cpp_obj = test_cpp.create_obj()

    print("Object ID at start: " + str(id(cpp_obj)))

    # generate new string and add
    new_string = [options.str] * options.num
    cpp_obj.add(new_string)

    print("Post-construction: " + str(p0.memory_percent()))

    # create new set of additions
    additions = [str(i) for i in range(0, options.iter)]

    # crate shared memory object
    total_arr = np.array([cpp_obj])
    with SharedMemoryManager() as smm:
        array_shd, array_shd_tup = generate_shared_mem_array(total_arr, smm)

        # iterate multithreaded
        with Pool(processes=options.threads, maxtasksperchild=options.maxtasks) as pool:
            for _ in pool.imap(partial(addition_wrapper, array_shd_tup=array_shd_tup),
                                                  enumerate(additions)):
                pass

    print("Post-iteration: " + str(p0.memory_percent()))


if __name__ == '__main__':
    main()

shared_memory.py

import sys
import numpy as np
import collections

try:
    from multiprocessing import Pool, shared_memory
    from multiprocessing.managers import SharedMemoryManager

    NumpyShared = collections.namedtuple('NumpyShared', ('name', 'shape', 'dtype'))
except ImportError as e:
    sys.stderr.write("This version of ggCaller requires python v3.8 or higher\n")
    sys.exit(1)


# generate shared memory array
def generate_shared_mem_array(in_array, smm):
    array_raw = smm.SharedMemory(size=in_array.nbytes)
    array_shared = np.ndarray(in_array.shape, dtype=in_array.dtype, buffer=array_raw.buf)
    array_shared[:] = in_array[:]
    array_shared_tup = NumpyShared(name=array_raw.name, shape=in_array.shape, dtype=in_array.dtype)
    return (array_shared, array_shared_tup)
0 Answers
Related