mirror of
https://github.com/wassname/ray.git
synced 2026-09-12 12:51:15 +08:00
Availability after worker failure (#316)
* Availability after a killed worker * Workers exit cleanly * Memory cleanup in photon C tests * Worker failure in multinode * Consolidate worker cleanup handlers * Update the result table before handling a task submission * KILL_WORKER_TIMEOUT -> KILL_WORKER_TIMEOUT_MILLISECONDS * Log a warning instead of crashing if no result table entry found
This commit is contained in:
committed by
Robert Nishihara
parent
232601f90d
commit
be1618f041
@@ -70,5 +70,38 @@ class ComponentFailureTest(unittest.TestCase):
|
||||
self.assertTrue(ray.services.all_processes_alive(exclude=[ray.services.PROCESS_TYPE_WORKER]))
|
||||
ray.worker.cleanup()
|
||||
|
||||
def _testWorkerFailed(self, num_local_schedulers):
|
||||
@ray.remote
|
||||
def f(x):
|
||||
time.sleep(0.5)
|
||||
return x
|
||||
|
||||
num_initial_workers = 4
|
||||
ray.worker._init(num_workers=num_initial_workers * num_local_schedulers,
|
||||
num_local_schedulers=num_local_schedulers,
|
||||
start_workers_from_local_scheduler=False,
|
||||
start_ray_local=True,
|
||||
num_cpus=[num_initial_workers] * num_local_schedulers)
|
||||
# Submit more tasks than there are workers so that all workers and cores
|
||||
# are utilized.
|
||||
object_ids = [f.remote(i) for i in range(num_initial_workers * num_local_schedulers)]
|
||||
object_ids += [f.remote(object_id) for object_id in object_ids]
|
||||
# Allow the tasks some time to begin executing.
|
||||
time.sleep(0.1)
|
||||
# Kill the workers as the tasks execute.
|
||||
for worker in ray.services.all_processes[ray.services.PROCESS_TYPE_WORKER]:
|
||||
worker.terminate()
|
||||
time.sleep(0.1)
|
||||
# Make sure that we can still get the objects after the executing tasks died.
|
||||
ray.get(object_ids)
|
||||
|
||||
ray.worker.cleanup()
|
||||
|
||||
def testWorkerFailed(self):
|
||||
self._testWorkerFailed(1)
|
||||
|
||||
def testWorkerFailedMultinode(self):
|
||||
self._testWorkerFailed(4)
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main(verbosity=2)
|
||||
|
||||
Reference in New Issue
Block a user