From b9e1977fae6c44db125918ed47208e1919b6da97 Mon Sep 17 00:00:00 2001 From: Yuhong Guo Date: Fri, 7 Dec 2018 04:55:49 +0800 Subject: [PATCH] Fix failure of test_free_objects_multi_node (#3481) It is possible that `test_free_objects_multi_node` would fail sometimes. If we run this test 20 times, we may found at least one failure. The cause is that the test is based on function tasks. One raylet may create more than one worker to execute the tasks. So flush operations may be separated to several workers and not clean all the worker objects held by the plasma client. In this PR, I change function task to actor tasks, which guarantee all the tasks are executed in one worker of a raylet. --- test/runtest.py | 66 ++++++++++++++++++++++++++++++++----------------- 1 file changed, 44 insertions(+), 22 deletions(-) diff --git a/test/runtest.py b/test/runtest.py index 22251efd7..91862023b 100644 --- a/test/runtest.py +++ b/test/runtest.py @@ -3,6 +3,7 @@ from __future__ import division from __future__ import print_function import json +import logging import os import re import setproctitle @@ -21,6 +22,8 @@ import ray.ray_constants as ray_constants import ray.test.cluster_utils import ray.test.test_utils +logger = logging.getLogger(__name__) + def assert_equal(obj1, obj2): module_numpy = (type(obj1).__module__ == np.__name__ @@ -699,7 +702,7 @@ def test_defining_remote_functions(shutdown_only): if val == 10: break else: - print("Still using old definition of f, trying again.") + logger.info("Still using old definition of f, trying again.") # Test that we can close over plain old data. data = [ @@ -1248,6 +1251,17 @@ def test_multithreading(shutdown_only): def test_free_objects_multi_node(shutdown_only): + # This test will do following: + # 1. Create 3 raylets that each hold an actor. + # 2. Each actor creates an object which is the deletion target. + # 3. Invoke 64 methods on each actor to flush plasma client. + # 4. After flushing, the plasma client releases the targets. + # 5. Check that the deletion targets have been deleted. + # Caution: if remote functions are used instead of actor methods, + # one raylet may create more than one worker to execute the + # tasks, so the flushing operations may be executed in different + # workers and the plasma client holding the deletion target + # may not be flushed. config = json.dumps({"object_manager_repeated_push_delay_ms": 1000}) ray.worker._init( start_ray_local=True, @@ -1263,53 +1277,61 @@ def test_free_objects_multi_node(shutdown_only): _internal_config=config) @ray.remote(resources={"Custom0": 1}) - def run_on_0(): - return ray.worker.global_worker.plasma_client.store_socket_name + class ActorOnNode0(object): + def get(self): + return ray.worker.global_worker.plasma_client.store_socket_name @ray.remote(resources={"Custom1": 1}) - def run_on_1(): - return ray.worker.global_worker.plasma_client.store_socket_name + class ActorOnNode1(object): + def get(self): + return ray.worker.global_worker.plasma_client.store_socket_name @ray.remote(resources={"Custom2": 1}) - def run_on_2(): - return ray.worker.global_worker.plasma_client.store_socket_name + class ActorOnNode2(object): + def get(self): + return ray.worker.global_worker.plasma_client.store_socket_name - def create(): - a = run_on_0.remote() - b = run_on_1.remote() - c = run_on_2.remote() + def create(actors): + a = actors[0].get.remote() + b = actors[1].get.remote() + c = actors[2].get.remote() (l1, l2) = ray.wait([a, b, c], num_returns=3) assert len(l1) == 3 assert len(l2) == 0 return (a, b, c) - def flush(): + def flush(actors): # Flush the Release History. # Current Plasma Client Cache will maintain 64-item list. # If the number changed, this will fail. - print("Start Flush!") + logger.info("Start Flush!") for i in range(64): - ray.get([run_on_0.remote(), run_on_1.remote(), run_on_2.remote()]) - print("Flush finished!") + ray.get([actor.get.remote() for actor in actors]) + logger.info("Flush finished!") - def run_one_test(local_only): - (a, b, c) = create() + def run_one_test(actors, local_only): + (a, b, c) = create(actors) # The three objects should be generated on different object stores. assert ray.get(a) != ray.get(b) assert ray.get(a) != ray.get(c) assert ray.get(c) != ray.get(b) ray.internal.free([a, b, c], local_only=local_only) - flush() + flush(actors) return (a, b, c) + actors = [ + ActorOnNode0.remote(), + ActorOnNode1.remote(), + ActorOnNode2.remote() + ] # Case 1: run this local_only=False. All 3 objects will be deleted. - (a, b, c) = run_one_test(False) + (a, b, c) = run_one_test(actors, False) (l1, l2) = ray.wait([a, b, c], timeout=10, num_returns=1) # All the objects are deleted. assert len(l1) == 0 assert len(l2) == 3 # Case 2: run this local_only=True. Only 1 object will be deleted. - (a, b, c) = run_one_test(True) + (a, b, c) = run_one_test(actors, True) (l1, l2) = ray.wait([a, b, c], timeout=10, num_returns=3) # One object is deleted and 2 objects are not. assert len(l1) == 2 @@ -2114,7 +2136,7 @@ def attempt_to_load_balance(remote_function, [remote_function.remote(*args) for _ in range(total_tasks)]) names = set(locations) counts = [locations.count(name) for name in names] - print("Counts are {}.".format(counts)) + logger.info("Counts are {}.".format(counts)) if (len(names) == num_local_schedulers and all(count >= minimum_count for count in counts)): break @@ -2299,7 +2321,7 @@ def test_log_file_api(shutdown_only): @ray.remote def f(): - print(message) + logger.info(message) # The call to sys.stdout.flush() seems to be necessary when using # the system Python 2.7 on Ubuntu. sys.stdout.flush()