| 27 | namespace oneflow { |
| 28 | |
| 29 | Thread::Thread(const StreamId& stream_id) : thrd_id_(EncodeStreamIdToInt64(stream_id)) { |
| 30 | local_msg_queue_enabled_ = ParseBooleanFromEnv("ONEFLOW_THREAD_ENABLE_LOCAL_MESSAGE_QUEUE", true); |
| 31 | light_actor_enabled_ = ParseBooleanFromEnv("ONEFLOW_ACTOR_ENABLE_LIGHT_ACTOR", true); |
| 32 | if (IsClassRegistered<int, StreamContext, const StreamId&>(stream_id.device_id().device_type(), |
| 33 | stream_id)) { |
| 34 | stream_ctx_.reset(NewObj<int, StreamContext, const StreamId&>( |
| 35 | stream_id.device_id().device_type(), stream_id)); |
| 36 | } else { |
| 37 | stream_ctx_.reset(new GenericStreamContext(stream_id)); |
| 38 | } |
| 39 | |
| 40 | actor_thread_ = std::thread([this, stream_id]() { |
| 41 | LazyMode::Guard guard(true); |
| 42 | OF_PROFILER_NAME_THIS_HOST_THREAD("_" + ToString(stream_id.device_id().device_type()) |
| 43 | + std::to_string(stream_id.device_id().device_index()) |
| 44 | + "_actor"); |
| 45 | CHECK_JUST(stream_ctx_->stream()->OnExecutionContextSetup()); |
| 46 | PollMsgChannel(); |
| 47 | CHECK_JUST(stream_ctx_->stream()->OnExecutionContextTeardown()); |
| 48 | }); |
| 49 | } |
| 50 | |
| 51 | Thread::~Thread() { |
| 52 | actor_thread_.join(); |