mirror of
https://github.com/wassname/ray.git
synced 2026-07-20 12:40:20 +08:00
Add test to test raylet client connection when raylet crashes. (#3518)
This commit is contained in:
committed by
Robert Nishihara
parent
e7b51cbd1b
commit
a4abe6c0fe
@@ -6,7 +6,6 @@ import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import org.ray.api.RayObject;
|
||||
import org.ray.api.WaitResult;
|
||||
import org.ray.api.exception.RayException;
|
||||
import org.ray.api.id.UniqueId;
|
||||
import org.ray.runtime.RayDevRuntime;
|
||||
import org.ray.runtime.objectstore.MockObjectStore;
|
||||
@@ -68,7 +67,7 @@ public class MockRayletClient implements RayletClient {
|
||||
|
||||
@Override
|
||||
public void fetchOrReconstruct(List<UniqueId> objectIds, boolean fetchOnly,
|
||||
UniqueId currentTaskId) throws RayException {
|
||||
UniqueId currentTaskId) {
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -3,7 +3,6 @@ package org.ray.runtime.raylet;
|
||||
import java.util.List;
|
||||
import org.ray.api.RayObject;
|
||||
import org.ray.api.WaitResult;
|
||||
import org.ray.api.exception.RayException;
|
||||
import org.ray.api.id.UniqueId;
|
||||
import org.ray.runtime.task.TaskSpec;
|
||||
|
||||
@@ -16,8 +15,7 @@ public interface RayletClient {
|
||||
|
||||
TaskSpec getTask();
|
||||
|
||||
void fetchOrReconstruct(List<UniqueId> objectIds, boolean fetchOnly, UniqueId currentTaskId)
|
||||
throws RayException;
|
||||
void fetchOrReconstruct(List<UniqueId> objectIds, boolean fetchOnly, UniqueId currentTaskId);
|
||||
|
||||
void notifyUnblocked(UniqueId currentTaskId);
|
||||
|
||||
|
||||
@@ -90,7 +90,7 @@ public class RayletClientImpl implements RayletClient {
|
||||
|
||||
@Override
|
||||
public void fetchOrReconstruct(List<UniqueId> objectIds, boolean fetchOnly,
|
||||
UniqueId currentTaskId) throws RayException {
|
||||
UniqueId currentTaskId) {
|
||||
if (RayLog.core.isDebugEnabled()) {
|
||||
RayLog.core.debug("Blocked on objects for task {}, object IDs are {}",
|
||||
UniqueIdUtil.computeTaskId(objectIds.get(0)), objectIds);
|
||||
|
||||
@@ -98,7 +98,7 @@ static PyObject *PyRayletClient_FetchOrReconstruct(PyRayletClient *self, PyObjec
|
||||
<< "raylet client may be closed, check raylet status. error message: "
|
||||
<< status.ToString();
|
||||
PyErr_SetString(CommonError, stream.str().c_str());
|
||||
Py_RETURN_NONE;
|
||||
return NULL;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -8,9 +8,11 @@ import os
|
||||
import ray
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
|
||||
import ray.ray_constants as ray_constants
|
||||
from ray.utils import _random_string
|
||||
import pytest
|
||||
|
||||
|
||||
@@ -611,3 +613,18 @@ def test_warning_for_dead_node(ray_start_two_nodes):
|
||||
}
|
||||
|
||||
assert client_ids == warning_client_ids
|
||||
|
||||
|
||||
def test_raylet_crash_when_get(ray_start_regular):
|
||||
nonexistent_id = ray.ObjectID(_random_string())
|
||||
|
||||
def sleep_to_kill_raylet():
|
||||
# Don't kill raylet before default workers get connected.
|
||||
time.sleep(2)
|
||||
ray.services.all_processes[ray.services.PROCESS_TYPE_RAYLET][0].kill()
|
||||
|
||||
thread = threading.Thread(target=sleep_to_kill_raylet)
|
||||
thread.start()
|
||||
with pytest.raises(Exception, match=r".*raylet client may be closed.*"):
|
||||
ray.get(nonexistent_id)
|
||||
thread.join()
|
||||
|
||||
Reference in New Issue
Block a user