mirror of
https://github.com/wassname/ray.git
synced 2026-09-12 12:51:15 +08:00
Submit task asynchronously from raylet client (#5313)
This commit is contained in:
@@ -2,6 +2,7 @@ from __future__ import absolute_import
|
|||||||
from __future__ import division
|
from __future__ import division
|
||||||
from __future__ import print_function
|
from __future__ import print_function
|
||||||
|
|
||||||
|
import json
|
||||||
import os
|
import os
|
||||||
import signal
|
import signal
|
||||||
import sys
|
import sys
|
||||||
@@ -311,8 +312,12 @@ def check_components_alive(cluster, component_type, check_component_alive):
|
|||||||
@pytest.mark.parametrize(
|
@pytest.mark.parametrize(
|
||||||
"ray_start_cluster", [{
|
"ray_start_cluster", [{
|
||||||
"num_cpus": 8,
|
"num_cpus": 8,
|
||||||
"num_nodes": 4
|
"num_nodes": 4,
|
||||||
}], indirect=True)
|
"_internal_config": json.dumps({
|
||||||
|
"num_heartbeats_timeout": 100
|
||||||
|
}),
|
||||||
|
}],
|
||||||
|
indirect=True)
|
||||||
def test_raylet_failed(ray_start_cluster):
|
def test_raylet_failed(ray_start_cluster):
|
||||||
cluster = ray_start_cluster
|
cluster = ray_start_cluster
|
||||||
# Kill all raylets on worker nodes.
|
# Kill all raylets on worker nodes.
|
||||||
@@ -329,8 +334,12 @@ def test_raylet_failed(ray_start_cluster):
|
|||||||
@pytest.mark.parametrize(
|
@pytest.mark.parametrize(
|
||||||
"ray_start_cluster", [{
|
"ray_start_cluster", [{
|
||||||
"num_cpus": 8,
|
"num_cpus": 8,
|
||||||
"num_nodes": 2
|
"num_nodes": 2,
|
||||||
}], indirect=True)
|
"_internal_config": json.dumps({
|
||||||
|
"num_heartbeats_timeout": 100
|
||||||
|
}),
|
||||||
|
}],
|
||||||
|
indirect=True)
|
||||||
def test_plasma_store_failed(ray_start_cluster):
|
def test_plasma_store_failed(ray_start_cluster):
|
||||||
cluster = ray_start_cluster
|
cluster = ray_start_cluster
|
||||||
# Kill all plasma stores on worker nodes.
|
# Kill all plasma stores on worker nodes.
|
||||||
|
|||||||
@@ -77,10 +77,18 @@ ray::Status RayletClient::SubmitTask(const ray::TaskSpecification &task_spec) {
|
|||||||
SubmitTaskRequest submit_task_request;
|
SubmitTaskRequest submit_task_request;
|
||||||
submit_task_request.mutable_task_spec()->CopyFrom(task_spec.GetMessage());
|
submit_task_request.mutable_task_spec()->CopyFrom(task_spec.GetMessage());
|
||||||
|
|
||||||
grpc::ClientContext context;
|
auto callback = [this](const Status &status, const SubmitTaskReply &reply) {
|
||||||
SubmitTaskReply reply;
|
if (!status.ok() && is_connected_) {
|
||||||
auto status = stub_->SubmitTask(&context, submit_task_request, &reply);
|
is_connected_ = false;
|
||||||
return GrpcStatusToRayStatus(status);
|
RAY_LOG(INFO) << "Failed to send SubmitTaskRequest, msg: " << status.message();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
auto call =
|
||||||
|
client_call_manager_.CreateCall<RayletService, SubmitTaskRequest, SubmitTaskReply>(
|
||||||
|
*stub_, &RayletService::Stub::PrepareAsyncSubmitTask, submit_task_request,
|
||||||
|
callback);
|
||||||
|
return call->GetStatus();
|
||||||
}
|
}
|
||||||
|
|
||||||
ray::Status RayletClient::GetTask(std::unique_ptr<ray::TaskSpecification> *task_spec) {
|
ray::Status RayletClient::GetTask(std::unique_ptr<ray::TaskSpecification> *task_spec) {
|
||||||
|
|||||||
Reference in New Issue
Block a user