mirror of
https://github.com/wassname/ray.git
synced 2026-08-13 12:30:18 +08:00
Support multiple core workers in one process (#7623)
This commit is contained in:
@@ -14,16 +14,14 @@ namespace streaming {
|
||||
class ReaderClient {
|
||||
public:
|
||||
/// Construct a ReaderClient object.
|
||||
/// \param[in] core_worker CoreWorker C++ pointer of current actor
|
||||
/// \param[in] async_func DataReader's raycall function descriptor to be called by
|
||||
/// DataWriter, asynchronous semantics \param[in] sync_func DataReader's raycall
|
||||
/// function descriptor to be called by DataWriter, synchronous semantics
|
||||
ReaderClient(CoreWorker *core_worker, RayFunction &async_func, RayFunction &sync_func)
|
||||
: core_worker_(core_worker) {
|
||||
ReaderClient(RayFunction &async_func, RayFunction &sync_func) {
|
||||
DownstreamQueueMessageHandler::peer_async_function_ = async_func;
|
||||
DownstreamQueueMessageHandler::peer_sync_function_ = sync_func;
|
||||
downstream_handler_ = ray::streaming::DownstreamQueueMessageHandler::CreateService(
|
||||
core_worker_, core_worker_->GetWorkerContext().GetCurrentActorID());
|
||||
CoreWorkerProcess::GetCoreWorker().GetWorkerContext().GetCurrentActorID());
|
||||
}
|
||||
|
||||
/// Post buffer to downstream queue service, asynchronously.
|
||||
@@ -34,19 +32,17 @@ class ReaderClient {
|
||||
std::shared_ptr<LocalMemoryBuffer> buffer);
|
||||
|
||||
private:
|
||||
CoreWorker *core_worker_;
|
||||
std::shared_ptr<DownstreamQueueMessageHandler> downstream_handler_;
|
||||
};
|
||||
|
||||
/// Interface of streaming queue for DataWriter. Similar to ReaderClient.
|
||||
class WriterClient {
|
||||
public:
|
||||
WriterClient(CoreWorker *core_worker, RayFunction &async_func, RayFunction &sync_func)
|
||||
: core_worker_(core_worker) {
|
||||
WriterClient(RayFunction &async_func, RayFunction &sync_func) {
|
||||
UpstreamQueueMessageHandler::peer_async_function_ = async_func;
|
||||
UpstreamQueueMessageHandler::peer_sync_function_ = sync_func;
|
||||
upstream_handler_ = ray::streaming::UpstreamQueueMessageHandler::CreateService(
|
||||
core_worker, core_worker_->GetWorkerContext().GetCurrentActorID());
|
||||
CoreWorkerProcess::GetCoreWorker().GetWorkerContext().GetCurrentActorID());
|
||||
}
|
||||
|
||||
void OnWriterMessage(std::shared_ptr<LocalMemoryBuffer> buffer);
|
||||
@@ -54,7 +50,6 @@ class WriterClient {
|
||||
std::shared_ptr<LocalMemoryBuffer> buffer);
|
||||
|
||||
private:
|
||||
CoreWorker *core_worker_;
|
||||
std::shared_ptr<UpstreamQueueMessageHandler> upstream_handler_;
|
||||
};
|
||||
} // namespace streaming
|
||||
|
||||
@@ -85,8 +85,8 @@ std::shared_ptr<Transport> QueueMessageHandler::GetOutTransport(
|
||||
void QueueMessageHandler::SetPeerActorID(const ObjectID &queue_id,
|
||||
const ActorID &actor_id) {
|
||||
actors_.emplace(queue_id, actor_id);
|
||||
out_transports_.emplace(
|
||||
queue_id, std::make_shared<ray::streaming::Transport>(core_worker_, actor_id));
|
||||
out_transports_.emplace(queue_id,
|
||||
std::make_shared<ray::streaming::Transport>(actor_id));
|
||||
}
|
||||
|
||||
ActorID QueueMessageHandler::GetPeerActorID(const ObjectID &queue_id) {
|
||||
@@ -113,10 +113,9 @@ void QueueMessageHandler::Stop() {
|
||||
}
|
||||
|
||||
std::shared_ptr<UpstreamQueueMessageHandler> UpstreamQueueMessageHandler::CreateService(
|
||||
CoreWorker *core_worker, const ActorID &actor_id) {
|
||||
const ActorID &actor_id) {
|
||||
if (nullptr == upstream_handler_) {
|
||||
upstream_handler_ =
|
||||
std::make_shared<UpstreamQueueMessageHandler>(core_worker, actor_id);
|
||||
upstream_handler_ = std::make_shared<UpstreamQueueMessageHandler>(actor_id);
|
||||
}
|
||||
return upstream_handler_;
|
||||
}
|
||||
@@ -247,11 +246,9 @@ void UpstreamQueueMessageHandler::ReleaseAllUpQueues() {
|
||||
}
|
||||
|
||||
std::shared_ptr<DownstreamQueueMessageHandler>
|
||||
DownstreamQueueMessageHandler::CreateService(CoreWorker *core_worker,
|
||||
const ActorID &actor_id) {
|
||||
DownstreamQueueMessageHandler::CreateService(const ActorID &actor_id) {
|
||||
if (nullptr == downstream_handler_) {
|
||||
downstream_handler_ =
|
||||
std::make_shared<DownstreamQueueMessageHandler>(core_worker, actor_id);
|
||||
downstream_handler_ = std::make_shared<DownstreamQueueMessageHandler>(actor_id);
|
||||
}
|
||||
return downstream_handler_;
|
||||
}
|
||||
|
||||
@@ -24,16 +24,9 @@ namespace streaming {
|
||||
class QueueMessageHandler {
|
||||
public:
|
||||
/// Construct a QueueMessageHandler instance.
|
||||
/// \param[in] core_worker CoreWorker C++ pointer of current actor, used to call Core
|
||||
/// Worker's api.
|
||||
/// For Python worker, the pointer can be obtained from
|
||||
/// ray.worker.global_worker.core_worker; For Java worker, obtained from
|
||||
/// RayNativeRuntime object through java reflection.
|
||||
/// \param[in] actor_id actor id of current actor.
|
||||
QueueMessageHandler(CoreWorker *core_worker, const ActorID &actor_id)
|
||||
: core_worker_(core_worker),
|
||||
actor_id_(actor_id),
|
||||
queue_dummy_work_(queue_service_) {
|
||||
QueueMessageHandler(const ActorID &actor_id)
|
||||
: actor_id_(actor_id), queue_dummy_work_(queue_service_) {
|
||||
Start();
|
||||
}
|
||||
|
||||
@@ -87,8 +80,6 @@ class QueueMessageHandler {
|
||||
void QueueThreadCallback() { queue_service_.run(); }
|
||||
|
||||
protected:
|
||||
/// CoreWorker C++ pointer of current actor
|
||||
CoreWorker *core_worker_;
|
||||
/// actor_id actor id of current actor
|
||||
ActorID actor_id_;
|
||||
/// Helper function, parse message buffer to Message object.
|
||||
@@ -111,8 +102,7 @@ class QueueMessageHandler {
|
||||
class UpstreamQueueMessageHandler : public QueueMessageHandler {
|
||||
public:
|
||||
/// Construct a UpstreamQueueMessageHandler instance.
|
||||
UpstreamQueueMessageHandler(CoreWorker *core_worker, const ActorID &actor_id)
|
||||
: QueueMessageHandler(core_worker, actor_id) {}
|
||||
UpstreamQueueMessageHandler(const ActorID &actor_id) : QueueMessageHandler(actor_id) {}
|
||||
/// Create a upstream queue.
|
||||
/// \param[in] queue_id queue id of the queue to be created.
|
||||
/// \param[in] peer_actor_id actor id of peer actor.
|
||||
@@ -140,7 +130,7 @@ class UpstreamQueueMessageHandler : public QueueMessageHandler {
|
||||
std::function<void(std::shared_ptr<LocalMemoryBuffer>)> callback) override;
|
||||
|
||||
static std::shared_ptr<UpstreamQueueMessageHandler> CreateService(
|
||||
CoreWorker *core_worker, const ActorID &actor_id);
|
||||
const ActorID &actor_id);
|
||||
static std::shared_ptr<UpstreamQueueMessageHandler> GetService();
|
||||
|
||||
static RayFunction peer_sync_function_;
|
||||
@@ -157,8 +147,8 @@ class UpstreamQueueMessageHandler : public QueueMessageHandler {
|
||||
/// UpstreamQueueMessageHandler holds and manages all downstream queues of current actor.
|
||||
class DownstreamQueueMessageHandler : public QueueMessageHandler {
|
||||
public:
|
||||
DownstreamQueueMessageHandler(CoreWorker *core_worker, const ActorID &actor_id)
|
||||
: QueueMessageHandler(core_worker, actor_id) {}
|
||||
DownstreamQueueMessageHandler(const ActorID &actor_id)
|
||||
: QueueMessageHandler(actor_id) {}
|
||||
std::shared_ptr<ReaderQueue> CreateDownstreamQueue(const ObjectID &queue_id,
|
||||
const ActorID &peer_actor_id);
|
||||
bool DownstreamQueueExists(const ObjectID &queue_id);
|
||||
@@ -178,7 +168,7 @@ class DownstreamQueueMessageHandler : public QueueMessageHandler {
|
||||
std::function<void(std::shared_ptr<LocalMemoryBuffer>)> callback);
|
||||
|
||||
static std::shared_ptr<DownstreamQueueMessageHandler> CreateService(
|
||||
CoreWorker *core_worker, const ActorID &actor_id);
|
||||
const ActorID &actor_id);
|
||||
static std::shared_ptr<DownstreamQueueMessageHandler> GetService();
|
||||
static RayFunction peer_sync_function_;
|
||||
static RayFunction peer_async_function_;
|
||||
|
||||
@@ -28,10 +28,9 @@ void Transport::SendInternal(std::shared_ptr<LocalMemoryBuffer> buffer,
|
||||
args.emplace_back(TaskArg::PassByValue(std::make_shared<RayObject>(
|
||||
std::move(buffer), meta, std::vector<ObjectID>(), true)));
|
||||
|
||||
STREAMING_CHECK(core_worker_ != nullptr);
|
||||
std::vector<std::shared_ptr<RayObject>> results;
|
||||
ray::Status st =
|
||||
core_worker_->SubmitActorTask(peer_actor_id_, function, args, options, &return_ids);
|
||||
ray::Status st = CoreWorkerProcess::GetCoreWorker().SubmitActorTask(
|
||||
peer_actor_id_, function, args, options, &return_ids);
|
||||
if (!st.ok()) {
|
||||
STREAMING_LOG(ERROR) << "SubmitActorTask failed. " << st;
|
||||
}
|
||||
@@ -50,7 +49,8 @@ std::shared_ptr<LocalMemoryBuffer> Transport::SendForResult(
|
||||
SendInternal(buffer, function, TASK_OPTION_RETURN_NUM_1, return_ids);
|
||||
|
||||
std::vector<std::shared_ptr<RayObject>> results;
|
||||
Status get_st = core_worker_->Get(return_ids, timeout_ms, &results);
|
||||
Status get_st =
|
||||
CoreWorkerProcess::GetCoreWorker().Get(return_ids, timeout_ms, &results);
|
||||
if (!get_st.ok()) {
|
||||
STREAMING_LOG(ERROR) << "Get fail.";
|
||||
return nullptr;
|
||||
|
||||
@@ -13,11 +13,10 @@ namespace streaming {
|
||||
class Transport {
|
||||
public:
|
||||
/// Construct a Transport object.
|
||||
/// \param[in] core_worker CoreWorker C++ pointer of current actor, which we call direct
|
||||
/// actor call interface with.
|
||||
/// \param[in] peer_actor_id actor id of peer actor.
|
||||
Transport(CoreWorker *core_worker, const ActorID &peer_actor_id)
|
||||
: core_worker_(core_worker), peer_actor_id_(peer_actor_id) {}
|
||||
Transport(const ActorID &peer_actor_id)
|
||||
: worker_id_(CoreWorkerProcess::GetCoreWorker().GetWorkerID()),
|
||||
peer_actor_id_(peer_actor_id) {}
|
||||
virtual ~Transport() = default;
|
||||
|
||||
/// Send buffer asynchronously, peer's `function` will be called.
|
||||
@@ -55,7 +54,7 @@ class Transport {
|
||||
std::vector<ObjectID> &return_ids);
|
||||
|
||||
private:
|
||||
CoreWorker *core_worker_;
|
||||
WorkerID worker_id_;
|
||||
ActorID peer_actor_id_;
|
||||
};
|
||||
} // namespace streaming
|
||||
|
||||
Reference in New Issue
Block a user