From 6993a471f1ef781c93962eb9b43f41da21e51add Mon Sep 17 00:00:00 2001 From: Tianyi Chen Date: Fri, 13 Mar 2020 12:24:37 +0800 Subject: [PATCH] [Streaming] Move resource-manager and scheduler to master package. (#7582) --- .../master/resourcemanager/ResourceManager.java | 4 ++-- .../resourcemanager/ResourceManagerImpl.java | 6 +++--- .../scheduler/strategy/SlotAssignStrategy.java | 3 +-- .../strategy/SlotAssignStrategyFactory.java | 4 ++-- .../strategy/impl/PipelineFirstStrategy.java | 4 ++-- .../resourcemanager/ResourceManagerTest.java | 8 ++++---- .../strategy/PipelineFirstStrategyTest.java | 16 ++-------------- 7 files changed, 16 insertions(+), 29 deletions(-) rename streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/{core => }/master/resourcemanager/ResourceManager.java (89%) rename streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/{core => }/master/resourcemanager/ResourceManagerImpl.java (96%) rename streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/{core => }/master/scheduler/strategy/SlotAssignStrategy.java (94%) rename streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/{core => }/master/scheduler/strategy/SlotAssignStrategyFactory.java (81%) rename streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/{core => }/master/scheduler/strategy/impl/PipelineFirstStrategy.java (98%) diff --git a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/resourcemanager/ResourceManager.java b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/resourcemanager/ResourceManager.java similarity index 89% rename from streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/resourcemanager/ResourceManager.java rename to streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/resourcemanager/ResourceManager.java index 6313780c6..f8c544d5e 100644 --- a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/resourcemanager/ResourceManager.java +++ b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/resourcemanager/ResourceManager.java @@ -1,10 +1,10 @@ -package org.ray.streaming.runtime.core.master.resourcemanager; +package org.ray.streaming.runtime.master.resourcemanager; import java.util.List; import java.util.Map; -import org.ray.streaming.runtime.core.master.scheduler.strategy.SlotAssignStrategy; import org.ray.streaming.runtime.core.resource.Container; import org.ray.streaming.runtime.core.resource.Resources; +import org.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy; /** * The resource manager is responsible for resource de-/allocation and monitoring ray cluster. diff --git a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/resourcemanager/ResourceManagerImpl.java b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/resourcemanager/ResourceManagerImpl.java similarity index 96% rename from streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/resourcemanager/ResourceManagerImpl.java rename to streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/resourcemanager/ResourceManagerImpl.java index 235a6a576..0bad2111f 100644 --- a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/resourcemanager/ResourceManagerImpl.java +++ b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/resourcemanager/ResourceManagerImpl.java @@ -1,4 +1,4 @@ -package org.ray.streaming.runtime.core.master.resourcemanager; +package org.ray.streaming.runtime.master.resourcemanager; import java.util.ArrayList; import java.util.HashMap; @@ -13,11 +13,11 @@ import org.ray.api.runtimecontext.NodeInfo; import org.ray.streaming.runtime.config.StreamingMasterConfig; import org.ray.streaming.runtime.config.master.ResourceConfig; import org.ray.streaming.runtime.config.types.SlotAssignStrategyType; -import org.ray.streaming.runtime.core.master.scheduler.strategy.SlotAssignStrategy; -import org.ray.streaming.runtime.core.master.scheduler.strategy.SlotAssignStrategyFactory; import org.ray.streaming.runtime.core.resource.Container; import org.ray.streaming.runtime.core.resource.Resources; import org.ray.streaming.runtime.master.JobRuntimeContext; +import org.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy; +import org.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategyFactory; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/SlotAssignStrategy.java b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/SlotAssignStrategy.java similarity index 94% rename from streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/SlotAssignStrategy.java rename to streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/SlotAssignStrategy.java index 2a80a62fa..14ecde275 100644 --- a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/SlotAssignStrategy.java +++ b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/SlotAssignStrategy.java @@ -1,8 +1,7 @@ -package org.ray.streaming.runtime.core.master.scheduler.strategy; +package org.ray.streaming.runtime.master.scheduler.strategy; import java.util.List; import java.util.Map; - import org.ray.streaming.runtime.core.graph.executiongraph.ExecutionGraph; import org.ray.streaming.runtime.core.resource.Container; import org.ray.streaming.runtime.core.resource.ContainerID; diff --git a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/SlotAssignStrategyFactory.java b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/SlotAssignStrategyFactory.java similarity index 81% rename from streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/SlotAssignStrategyFactory.java rename to streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/SlotAssignStrategyFactory.java index 4dc19b9a0..bbafa1804 100644 --- a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/SlotAssignStrategyFactory.java +++ b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/SlotAssignStrategyFactory.java @@ -1,7 +1,7 @@ -package org.ray.streaming.runtime.core.master.scheduler.strategy; +package org.ray.streaming.runtime.master.scheduler.strategy; import org.ray.streaming.runtime.config.types.SlotAssignStrategyType; -import org.ray.streaming.runtime.core.master.scheduler.strategy.impl.PipelineFirstStrategy; +import org.ray.streaming.runtime.master.scheduler.strategy.impl.PipelineFirstStrategy; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/impl/PipelineFirstStrategy.java b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/impl/PipelineFirstStrategy.java similarity index 98% rename from streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/impl/PipelineFirstStrategy.java rename to streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/impl/PipelineFirstStrategy.java index ee99be986..8867225b9 100644 --- a/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/core/master/scheduler/strategy/impl/PipelineFirstStrategy.java +++ b/streaming/java/streaming-runtime/src/main/java/org/ray/streaming/runtime/master/scheduler/strategy/impl/PipelineFirstStrategy.java @@ -1,4 +1,4 @@ -package org.ray.streaming.runtime.core.master.scheduler.strategy.impl; +package org.ray.streaming.runtime.master.scheduler.strategy.impl; import com.google.common.base.Preconditions; import java.util.HashMap; @@ -8,11 +8,11 @@ import org.ray.streaming.runtime.config.types.SlotAssignStrategyType; import org.ray.streaming.runtime.core.graph.executiongraph.ExecutionGraph; import org.ray.streaming.runtime.core.graph.executiongraph.ExecutionJobVertex; import org.ray.streaming.runtime.core.graph.executiongraph.ExecutionVertex; -import org.ray.streaming.runtime.core.master.scheduler.strategy.SlotAssignStrategy; import org.ray.streaming.runtime.core.resource.Container; import org.ray.streaming.runtime.core.resource.ContainerID; import org.ray.streaming.runtime.core.resource.Resources; import org.ray.streaming.runtime.core.resource.Slot; +import org.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/resourcemanager/ResourceManagerTest.java b/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/resourcemanager/ResourceManagerTest.java index 77bdef73c..cec53f056 100644 --- a/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/resourcemanager/ResourceManagerTest.java +++ b/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/resourcemanager/ResourceManagerTest.java @@ -13,10 +13,10 @@ import org.ray.streaming.runtime.config.StreamingConfig; import org.ray.streaming.runtime.config.global.CommonConfig; import org.ray.streaming.runtime.config.master.ResourceConfig; import org.ray.streaming.runtime.core.graph.executiongraph.ExecutionGraph; -import org.ray.streaming.runtime.core.master.resourcemanager.ResourceManager; -import org.ray.streaming.runtime.core.master.resourcemanager.ResourceManagerImpl; -import org.ray.streaming.runtime.core.master.scheduler.strategy.SlotAssignStrategy; -import org.ray.streaming.runtime.core.master.scheduler.strategy.impl.PipelineFirstStrategy; +import org.ray.streaming.runtime.master.resourcemanager.ResourceManager; +import org.ray.streaming.runtime.master.resourcemanager.ResourceManagerImpl; +import org.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy; +import org.ray.streaming.runtime.master.scheduler.strategy.impl.PipelineFirstStrategy; import org.ray.streaming.runtime.core.resource.Container; import org.ray.streaming.runtime.core.resource.ContainerID; import org.ray.streaming.runtime.core.resource.Slot; diff --git a/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/schedule/strategy/PipelineFirstStrategyTest.java b/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/schedule/strategy/PipelineFirstStrategyTest.java index 66b6cb66d..4764b6d6f 100644 --- a/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/schedule/strategy/PipelineFirstStrategyTest.java +++ b/streaming/java/streaming-runtime/src/test/java/org/ray/streaming/runtime/schedule/strategy/PipelineFirstStrategyTest.java @@ -7,27 +7,15 @@ import java.util.Map; import java.util.Map.Entry; -import com.google.common.collect.Lists; import org.aeonbits.owner.ConfigFactory; -import org.ray.api.RayActor; -import org.ray.api.id.ActorId; -import org.ray.api.id.ObjectId; import org.ray.api.id.UniqueId; -import org.ray.runtime.actor.LocalModeRayActor; -import org.ray.streaming.api.context.RuntimeContext; -import org.ray.streaming.api.context.StreamingContext; -import org.ray.streaming.api.stream.DataStream; -import org.ray.streaming.api.stream.DataStreamSink; -import org.ray.streaming.api.stream.DataStreamSource; import org.ray.streaming.jobgraph.JobGraph; -import org.ray.streaming.jobgraph.JobGraphBuilder; import org.ray.streaming.runtime.BaseUnitTest; import org.ray.streaming.runtime.config.StreamingConfig; -import org.ray.streaming.runtime.config.StreamingMasterConfig; import org.ray.streaming.runtime.config.master.ResourceConfig; import org.ray.streaming.runtime.core.graph.executiongraph.ExecutionGraph; -import org.ray.streaming.runtime.core.master.scheduler.strategy.SlotAssignStrategy; -import org.ray.streaming.runtime.core.master.scheduler.strategy.impl.PipelineFirstStrategy; +import org.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy; +import org.ray.streaming.runtime.master.scheduler.strategy.impl.PipelineFirstStrategy; import org.ray.streaming.runtime.core.resource.Container; import org.ray.streaming.runtime.core.resource.ContainerID; import org.ray.streaming.runtime.core.resource.Resources;