Merge pull request #137 from amplab/memoryfix

Preparation to deallocate objects properly from object stores
This commit is contained in:
Robert Nishihara
2016-06-21 13:47:24 -07:00
committed by GitHub
6 changed files with 80 additions and 7 deletions
+10 -2
View File
@@ -512,7 +512,12 @@ PyObject* get_arrow(PyObject* self, PyObject* args) {
if (!PyArg_ParseTuple(args, "O&O&", &PyObjectToWorker, &worker, &PyObjectToObjRef, &objref)) {
return NULL;
}
return (PyObject*) worker->get_arrow(objref);
SegmentId segmentid;
PyObject* value = worker->get_arrow(objref, segmentid);
PyObject* val_and_segmentid = PyList_New(2);
PyList_SetItem(val_and_segmentid, 0, value);
PyList_SetItem(val_and_segmentid, 1, PyInt_FromLong(segmentid));
return val_and_segmentid;
}
PyObject* is_arrow(PyObject* self, PyObject* args) {
@@ -748,7 +753,10 @@ PyObject* get_object(PyObject* self, PyObject* args) {
slice s = worker->get_object(objref);
Obj* obj = new Obj(); // TODO: Make sure this will get deleted
obj->ParseFromString(std::string(reinterpret_cast<char*>(s.data), s.len));
return PyCapsule_New(static_cast<void*>(obj), "obj", &ObjCapsule_Destructor);
PyObject* result = PyList_New(2);
PyList_SetItem(result, 0, PyCapsule_New(static_cast<void*>(obj), "obj", &ObjCapsule_Destructor));
PyList_SetItem(result, 1, PyInt_FromLong(s.segmentid));
return result;
}
PyObject* request_object(PyObject* self, PyObject* args) {
+5 -1
View File
@@ -93,6 +93,7 @@ slice Worker::get_object(ObjRef objref) {
slice slice;
slice.data = segmentpool_->get_address(result);
slice.len = result.size();
slice.segmentid = result.segmentid();
return slice;
}
@@ -165,7 +166,9 @@ PyObject* Worker::put_arrow(ObjRef objref, PyObject* value) {
Py_RETURN_NONE;
}
PyObject* Worker::get_arrow(ObjRef objref) {
// returns python list containing the value represented by objref and the
// segmentid in which the object is stored
PyObject* Worker::get_arrow(ObjRef objref, SegmentId& segmentid) {
RAY_CHECK(connected_, "Attempted to perform get_arrow but failed.");
ObjRequest request;
request.workerid = workerid_;
@@ -176,6 +179,7 @@ PyObject* Worker::get_arrow(ObjRef objref) {
receive_obj_queue_.receive(&result);
uint8_t* address = segmentpool_->get_address(result);
auto source = std::make_shared<BufferMemorySource>(address, result.size());
segmentid = result.segmentid();
PyObject* value;
CHECK_ARROW_STATUS(pynumbuf::ReadPythonObjectFrom(source.get(), result.metadata_offset(), &value), "error during ReadPythonObjectFrom: ");
return value;
+1 -1
View File
@@ -57,7 +57,7 @@ class Worker {
// stores an arrow object to the local object store
PyObject* put_arrow(ObjRef objref, PyObject* array);
// gets an arrow object from the local object store
PyObject* get_arrow(ObjRef objref);
PyObject* get_arrow(ObjRef objref, SegmentId& segmentid);
// determine if the object stored in objref is an arrow object // TODO(pcm): more general mechanism for this?
bool is_arrow(ObjRef objref);
// make `alias_objref` refer to the same object that `target_objref` refers to