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)