mirror of
https://github.com/wassname/ray.git
synced 2026-08-03 13:10:57 +08:00
Make putting large objects work. (#411)
* putting large objects * add more checks * support large objects * fix test * fix linting * upgrade to latest arrow version * check malloc return code * print mmap file sizes * printing * revert to dlmalloc * add prints * more prints * add printing * printing * fix * update * fix * update * print * initialization * temp * fix * update * fix linting * comment out object_store_full tests * fix test * fix test * evict objects if dlmalloc fails * fix stresstests * Fix linting. * Uncomment large-memory tests. * Increase memory for docker image for jenkins tests. * Reduce large memory tests. * Further reduce large memory tests.
This commit is contained in:
committed by
Robert Nishihara
parent
1e84747e13
commit
4043769ba2
@@ -40,7 +40,8 @@ class TestLocalSchedulerClient(unittest.TestCase):
|
||||
def setUp(self):
|
||||
# Start Plasma store.
|
||||
plasma_store_name, self.p1 = plasma.start_plasma_store()
|
||||
self.plasma_client = plasma.PlasmaClient(plasma_store_name)
|
||||
self.plasma_client = plasma.PlasmaClient(plasma_store_name,
|
||||
release_delay=0)
|
||||
# Start a local scheduler.
|
||||
scheduler_name, self.p2 = local_scheduler.start_local_scheduler(
|
||||
plasma_store_name, use_valgrind=USE_VALGRIND)
|
||||
@@ -169,19 +170,15 @@ class TestLocalSchedulerClient(unittest.TestCase):
|
||||
t.start()
|
||||
|
||||
# Make one of the dependencies available.
|
||||
self.plasma_client.create(object_id1.id(), 1)
|
||||
buf = self.plasma_client.create(object_id1.id(), 1)
|
||||
self.plasma_client.seal(object_id1.id())
|
||||
# Release the object.
|
||||
del buf
|
||||
# Check that the thread is still waiting for a task.
|
||||
time.sleep(0.1)
|
||||
self.assertTrue(t.is_alive())
|
||||
|
||||
# Force eviction of the first dependency.
|
||||
num_objects = 4
|
||||
object_size = plasma.DEFAULT_PLASMA_STORE_MEMORY // num_objects
|
||||
for i in range(num_objects + 1):
|
||||
object_id = random_object_id()
|
||||
self.plasma_client.create(object_id.id(), object_size)
|
||||
self.plasma_client.seal(object_id.id())
|
||||
self.plasma_client.evict(plasma.DEFAULT_PLASMA_STORE_MEMORY)
|
||||
# Check that the thread is still waiting for a task.
|
||||
time.sleep(0.1)
|
||||
self.assertTrue(t.is_alive())
|
||||
|
||||
@@ -165,58 +165,26 @@ class TestPlasmaClient(unittest.TestCase):
|
||||
|
||||
# Create a list to keep some of the buffers in scope.
|
||||
memory_buffers = []
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 8, 0)
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 5 * 10 ** 8, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 8. Make sure that we can't create an object of
|
||||
# size 10 ** 8 + 1, but we can create one of size 10 ** 8.
|
||||
assert_create_raises_plasma_full(self, 10 ** 8 + 1)
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 10 ** 8, 0)
|
||||
# Remaining space is 5 * 10 ** 8. Make sure that we can't create an object
|
||||
# of size 5 * 10 ** 8 + 1, but we can create one of size 2 * 10 ** 8.
|
||||
assert_create_raises_plasma_full(self, 5 * 10 ** 8 + 1)
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 2 * 10 ** 8, 0)
|
||||
del memory_buffer
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 10 ** 8, 0)
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 2 * 10 ** 8, 0)
|
||||
del memory_buffer
|
||||
assert_create_raises_plasma_full(self, 10 ** 8 + 1)
|
||||
assert_create_raises_plasma_full(self, 5 * 10 ** 8 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 7, 0)
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 2 * 10 ** 8, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 7.
|
||||
assert_create_raises_plasma_full(self, 10 ** 7 + 1)
|
||||
# Remaining space is 3 * 10 ** 8.
|
||||
assert_create_raises_plasma_full(self, 3 * 10 ** 8 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 6, 0)
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 10 ** 8, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 6.
|
||||
assert_create_raises_plasma_full(self, 10 ** 6 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 5, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 5.
|
||||
assert_create_raises_plasma_full(self, 10 ** 5 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 4, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 4.
|
||||
assert_create_raises_plasma_full(self, 10 ** 4 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 3, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 3.
|
||||
assert_create_raises_plasma_full(self, 10 ** 3 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 2, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 2.
|
||||
assert_create_raises_plasma_full(self, 10 ** 2 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 1, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 1.
|
||||
assert_create_raises_plasma_full(self, 10 ** 1 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 9 * 10 ** 0, 0)
|
||||
memory_buffers.append(memory_buffer)
|
||||
# Remaining space is 10 ** 0.
|
||||
assert_create_raises_plasma_full(self, 10 ** 0 + 1)
|
||||
|
||||
_, memory_buffer, _ = create_object(self.plasma_client, 1, 0)
|
||||
# Remaining space is 2 * 10 ** 8.
|
||||
assert_create_raises_plasma_full(self, 2 * 10 ** 8 + 1)
|
||||
|
||||
def test_contains(self):
|
||||
fake_object_ids = [random_object_id() for _ in range(100)]
|
||||
|
||||
Reference in New Issue
Block a user