mirror of
https://github.com/wassname/ray.git
synced 2026-08-04 13:14:14 +08:00
Use random string for worker c++ logfile. (#378)
This commit is contained in:
committed by
Philipp Moritz
parent
e94ed2fc97
commit
b29fc0c481
+8
-13
@@ -666,11 +666,18 @@ static PyObject* create_worker(PyObject* self, PyObject* args) {
|
||||
// scheduler will choose the object store address.
|
||||
const char* objstore_address;
|
||||
int mode;
|
||||
if (!PyArg_ParseTuple(args, "sssi", &node_ip_address, &scheduler_address, &objstore_address, &mode)) {
|
||||
const char* log_file_name;
|
||||
if (!PyArg_ParseTuple(args, "sssis", &node_ip_address, &scheduler_address, &objstore_address, &mode, &log_file_name)) {
|
||||
return NULL;
|
||||
}
|
||||
// Set the logging file.
|
||||
create_log_dir_or_die(log_file_name);
|
||||
global_ray_config.log_to_file = true;
|
||||
global_ray_config.logfile.open(log_file_name);
|
||||
// Create the worker.
|
||||
bool is_driver = (mode != Mode::WORKER_MODE);
|
||||
Worker* worker = new Worker(std::string(node_ip_address), std::string(scheduler_address), static_cast<Mode>(mode));
|
||||
// Register the worker.
|
||||
worker->register_worker(std::string(node_ip_address), std::string(objstore_address), is_driver);
|
||||
|
||||
PyObject* t = PyTuple_New(2);
|
||||
@@ -1023,17 +1030,6 @@ static PyObject* dump_computation_graph(PyObject* self, PyObject* args) {
|
||||
Py_RETURN_NONE;
|
||||
}
|
||||
|
||||
static PyObject* set_log_config(PyObject* self, PyObject* args) {
|
||||
const char* log_file_name;
|
||||
if (!PyArg_ParseTuple(args, "s", &log_file_name)) {
|
||||
return NULL;
|
||||
}
|
||||
create_log_dir_or_die(log_file_name);
|
||||
global_ray_config.log_to_file = true;
|
||||
global_ray_config.logfile.open(log_file_name);
|
||||
Py_RETURN_NONE;
|
||||
}
|
||||
|
||||
static PyObject* kill_workers(PyObject* self, PyObject* args) {
|
||||
Worker* worker;
|
||||
if (!PyArg_ParseTuple(args, "O&", &PyObjectToWorker, &worker)) {
|
||||
@@ -1074,7 +1070,6 @@ static PyMethodDef RayLibMethods[] = {
|
||||
{ "export_remote_function", export_remote_function, METH_VARARGS, "export a remote function to workers" },
|
||||
{ "export_reusable_variable", export_reusable_variable, METH_VARARGS, "export a reusable variable to the workers" },
|
||||
{ "dump_computation_graph", dump_computation_graph, METH_VARARGS, "dump the current computation graph to a file" },
|
||||
{ "set_log_config", set_log_config, METH_VARARGS, "set filename for raylib logging" },
|
||||
{ "kill_workers", kill_workers, METH_VARARGS, "kills all of the workers" },
|
||||
{ NULL, NULL, 0, NULL }
|
||||
};
|
||||
|
||||
+2
-2
@@ -1138,11 +1138,11 @@ int main(int argc, char** argv) {
|
||||
const char* scheduling_algorithm_name = get_cmd_option(argv, argv + argc, "--scheduler-algorithm");
|
||||
if (scheduling_algorithm_name) {
|
||||
if (std::string(scheduling_algorithm_name) == "naive") {
|
||||
RAY_LOG(RAY_INFO, "scheduler: using 'naive' scheduler" << std::endl);
|
||||
RAY_LOG(RAY_INFO, "scheduler: using 'naive' scheduler");
|
||||
scheduling_algorithm = SCHEDULING_ALGORITHM_NAIVE;
|
||||
}
|
||||
if (std::string(scheduling_algorithm_name) == "locality_aware") {
|
||||
RAY_LOG(RAY_INFO, "scheduler: using 'locality aware' scheduler" << std::endl);
|
||||
RAY_LOG(RAY_INFO, "scheduler: using 'locality aware' scheduler");
|
||||
scheduling_algorithm = SCHEDULING_ALGORITHM_LOCALITY_AWARE;
|
||||
}
|
||||
}
|
||||
|
||||
+5
-5
@@ -13,7 +13,7 @@ extern "C" {
|
||||
|
||||
inline WorkerServiceImpl::WorkerServiceImpl(const std::string& send_queue_name, Mode mode)
|
||||
: mode_(mode) {
|
||||
RAY_LOG(RAY_DEBUG, "Worker service connecting to queue " << send_queue_name);
|
||||
RAY_LOG(RAY_INFO, "Worker service connecting to queue " << send_queue_name);
|
||||
RAY_CHECK(send_queue_.connect(send_queue_name, false), "error connecting send_queue_");
|
||||
}
|
||||
|
||||
@@ -101,7 +101,7 @@ Worker::Worker(const std::string& node_ip_address, const std::string& scheduler_
|
||||
std::mt19937 rng(rd());
|
||||
std::uniform_int_distribution<int> queue_name_generator(0, 10000000);
|
||||
receive_queue_name_ = "worker_receive_queue:" + std::to_string(queue_name_generator(rng));
|
||||
RAY_LOG(RAY_DEBUG, "Worker creating queue " << receive_queue_name_ << std::endl);
|
||||
RAY_LOG(RAY_INFO, "Worker creating queue " << receive_queue_name_);
|
||||
RAY_CHECK(receive_queue_.connect(receive_queue_name_, true), "error connecting receive_queue_");
|
||||
}
|
||||
|
||||
@@ -162,11 +162,11 @@ void Worker::register_worker(const std::string& node_ip_address, const std::stri
|
||||
segmentpool_ = std::make_shared<MemorySegmentPool>(objstoreid_, objstore_address_, false);
|
||||
// Connect to the queue for sending requests to the object store.
|
||||
std::string request_obj_queue_name = std::string("queue:") + objstore_address_ + std::string(":obj");
|
||||
RAY_LOG(RAY_DEBUG, "Worker connecting to queue with name " << request_obj_queue_name << " to send requests to the object store.");
|
||||
RAY_LOG(RAY_INFO, "Worker connecting to queue with name " << request_obj_queue_name << " to send requests to the object store.");
|
||||
RAY_CHECK(request_obj_queue_.connect(request_obj_queue_name, false), "error connecting request_obj_queue_");
|
||||
// Create a queue for receiving messages from the object store.
|
||||
std::string receive_obj_queue_name = std::string("queue:") + objstore_address_ + std::string(":worker:") + std::to_string(workerid_) + std::string(":obj");
|
||||
RAY_LOG(RAY_DEBUG, "Worker creating queue with name " << receive_obj_queue_name << " to receive messages from the object store.");
|
||||
RAY_LOG(RAY_INFO, "Worker creating queue with name " << receive_obj_queue_name << " to receive messages from the object store.");
|
||||
RAY_CHECK(receive_obj_queue_.connect(receive_obj_queue_name, true), "error connecting receive_obj_queue_");
|
||||
connected_ = true;
|
||||
return;
|
||||
@@ -481,7 +481,7 @@ void Worker::start_worker_service(Mode mode) {
|
||||
// Wait for the worker service to start. This essentially implements a
|
||||
// condition variable using atomics, but that failed on Mac OS X on Travis.
|
||||
while (!worker_service_started.load()) {
|
||||
RAY_LOG(RAY_DEBUG, "Looping while waiting for the worker service to start.");
|
||||
RAY_LOG(RAY_INFO, "Looping while waiting for the worker service to start.");
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user