[streaming] Async changes for resourcemanager part (#7955)

This commit is contained in:
JianZhangYang
2020-04-16 14:15:45 +08:00
committed by GitHub
parent 5c274fe631
commit 7b0518b993
29 changed files with 907 additions and 842 deletions
@@ -36,7 +36,7 @@ public class ExecutionGraphTest extends BaseUnitTest {
GraphManager graphManager = new GraphManagerImpl(new JobRuntimeContext(streamingConfig));
JobGraph jobGraph = buildJobGraph();
ExecutionGraph executionGraph = buildExecutionGraph(graphManager, jobGraph);
List<ExecutionJobVertex> executionJobVertices = executionGraph.getExecutionJobVertexLices();
List<ExecutionJobVertex> executionJobVertices = executionGraph.getExecutionJobVertexList();
Assert.assertEquals(executionJobVertices.size(), jobGraph.getJobVertexList().size());
@@ -88,6 +88,7 @@ public class ExecutionGraphTest extends BaseUnitTest {
jobConfig.put("key1", "value1");
jobConfig.put("key2", "value2");
jobConfig.put(ResourceConfig.TASK_RESOURCE_CPU, "2.0");
jobConfig.put(ResourceConfig.TASK_RESOURCE_MEM, "2.0");
JobGraphBuilder jobGraphBuilder = new JobGraphBuilder(
Lists.newArrayList(streamSink), "test", jobConfig);
@@ -1,99 +1,80 @@
package io.ray.streaming.runtime.resourcemanager;
import io.ray.api.Ray;
import io.ray.streaming.jobgraph.JobGraph;
import io.ray.streaming.runtime.BaseUnitTest;
import io.ray.streaming.runtime.TestHelper;
import io.ray.streaming.runtime.config.StreamingConfig;
import io.ray.streaming.runtime.config.global.CommonConfig;
import io.ray.streaming.runtime.config.master.ResourceConfig;
import io.ray.streaming.runtime.core.graph.executiongraph.ExecutionGraph;
import io.ray.streaming.runtime.core.resource.Container;
import io.ray.streaming.runtime.core.resource.ContainerID;
import io.ray.streaming.runtime.core.resource.ResourceType;
import io.ray.streaming.runtime.core.resource.Slot;
import io.ray.streaming.runtime.graph.ExecutionGraphTest;
import io.ray.streaming.runtime.master.JobRuntimeContext;
import io.ray.streaming.runtime.master.graphmanager.GraphManager;
import io.ray.streaming.runtime.master.graphmanager.GraphManagerImpl;
import io.ray.streaming.runtime.master.resourcemanager.ResourceManager;
import io.ray.streaming.runtime.master.resourcemanager.ResourceManagerImpl;
import io.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy;
import io.ray.streaming.runtime.master.scheduler.strategy.impl.PipelineFirstStrategy;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import org.aeonbits.owner.util.Collections;
import org.mockito.MockitoAnnotations;
import org.powermock.core.classloader.annotations.PowerMockIgnore;
import org.powermock.core.classloader.annotations.PrepareForTest;
import io.ray.api.Ray;
import io.ray.api.id.UniqueId;
import io.ray.api.runtimecontext.NodeInfo;
import io.ray.streaming.runtime.config.StreamingConfig;
import io.ray.streaming.runtime.config.global.CommonConfig;
import io.ray.streaming.runtime.master.resourcemanager.ResourceManager;
import io.ray.streaming.runtime.master.resourcemanager.ResourceManagerImpl;
import io.ray.streaming.runtime.core.resource.Container;
import io.ray.streaming.runtime.master.JobRuntimeContext;
import io.ray.streaming.runtime.util.Mockitools;
import io.ray.streaming.runtime.util.RayUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.testng.Assert;
import org.testng.IObjectFactory;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.ObjectFactory;
import org.testng.annotations.Test;
public class ResourceManagerTest extends BaseUnitTest {
@PrepareForTest(RayUtils.class)
@PowerMockIgnore({"org.slf4j.*", "javax.xml.*"})
public class ResourceManagerTest {
private static final Logger LOG = LoggerFactory.getLogger(ResourceManagerTest.class);
private Object rayAsyncContext;
@ObjectFactory
public IObjectFactory getObjectFactory() {
return new org.powermock.modules.testng.PowerMockObjectFactory();
}
@org.testng.annotations.BeforeClass
public void setUp() {
// ray init
Ray.init();
TestHelper.setUTFlag();
LOG.warn("Do set up");
MockitoAnnotations.initMocks(this);
}
@org.testng.annotations.AfterClass
public void tearDown() {
TestHelper.clearUTFlag();
LOG.warn("Do tear down");
}
@BeforeMethod
public void mockGscApi() {
// ray init
Ray.init();
rayAsyncContext = Ray.getAsyncContext();
Mockitools.mockGscApi();
}
@Test
public void testGcsMockedApi() {
Map<UniqueId, NodeInfo> nodeInfoMap = RayUtils.getAliveNodeInfoMap();
Assert.assertEquals(nodeInfoMap.size(), 5);
}
@Test
public void testApi() {
Ray.setAsyncContext(rayAsyncContext);
Map<String, String> conf = new HashMap<String, String>();
conf.put(CommonConfig.JOB_NAME, "testApi");
conf.put(ResourceConfig.TASK_RESOURCE_CPU_LIMIT_ENABLE, "true");
conf.put(ResourceConfig.TASK_RESOURCE_MEM_LIMIT_ENABLE, "true");
conf.put(ResourceConfig.TASK_RESOURCE_MEM, "10");
conf.put(ResourceConfig.TASK_RESOURCE_CPU, "2");
StreamingConfig config = new StreamingConfig(conf);
JobRuntimeContext jobRuntimeContext = new JobRuntimeContext(config);
ResourceManager resourceManager = new ResourceManagerImpl(jobRuntimeContext);
SlotAssignStrategy slotAssignStrategy = resourceManager.getSlotAssignStrategy();
Assert.assertTrue(slotAssignStrategy instanceof PipelineFirstStrategy);
Map<String, Double> containerResource = new HashMap<>();
containerResource.put(ResourceType.CPU.name(), 16.0);
containerResource.put(ResourceType.MEM.name(), 128.0);
Container container1 = new Container(null, "testAddress1", "testHostName1");
container1.setAvailableResource(containerResource);
Container container2 = new Container(null, "testAddress2", "testHostName2");
container2.setAvailableResource(new HashMap<>(containerResource));
List<Container> containers = Collections.list(container1, container2);
resourceManager.getResources().getRegisterContainers().addAll(containers);
Assert.assertEquals(resourceManager.getRegisteredContainers().size(), 2);
//build ExecutionGraph
GraphManager graphManager = new GraphManagerImpl(new JobRuntimeContext(config));
JobGraph jobGraph = ExecutionGraphTest.buildJobGraph();
ExecutionGraph executionGraph = ExecutionGraphTest.buildExecutionGraph(graphManager, jobGraph);
int slotNumPerContainer = slotAssignStrategy.getSlotNumPerContainer(containers, executionGraph
.getMaxParallelism());
Assert.assertEquals(slotNumPerContainer, 1);
slotAssignStrategy.allocateSlot(containers, slotNumPerContainer);
Map<ContainerID, List<Slot>> allocatingMap = slotAssignStrategy.assignSlot(executionGraph);
Assert.assertEquals(allocatingMap.size(), 2);
executionGraph.getAllAddedExecutionVertices().forEach(vertex -> {
Container container = resourceManager.getResources()
.getRegisterContainerByContainerId(vertex.getSlot().getContainerID());
Map<String, Double> resource = resourceManager.allocateResource(container, vertex.getResources());
Assert.assertNotNull(resource);
});
Assert.assertEquals(container1.getAvailableResource().get(ResourceType.CPU.name()), 14.0);
Assert.assertEquals(container2.getAvailableResource().get(ResourceType.CPU.name()), 14.0);
Assert.assertEquals(container1.getAvailableResource().get(ResourceType.MEM.name()), 126.0);
Assert.assertEquals(container2.getAvailableResource().get(ResourceType.MEM.name()), 126.0);
// test register container
List<Container> containers = resourceManager.getRegisteredContainers();
Assert.assertEquals(containers.size(), 5);
}
}
@@ -4,28 +4,25 @@ import io.ray.api.id.UniqueId;
import io.ray.streaming.jobgraph.JobGraph;
import io.ray.streaming.runtime.BaseUnitTest;
import io.ray.streaming.runtime.config.StreamingConfig;
import io.ray.streaming.runtime.config.master.ResourceConfig;
import io.ray.streaming.runtime.config.types.ResourceAssignStrategyType;
import io.ray.streaming.runtime.core.graph.executiongraph.ExecutionGraph;
import io.ray.streaming.runtime.core.resource.Container;
import io.ray.streaming.runtime.core.resource.ContainerID;
import io.ray.streaming.runtime.core.resource.ResourceType;
import io.ray.streaming.runtime.core.resource.Resources;
import io.ray.streaming.runtime.core.resource.Slot;
import io.ray.streaming.runtime.graph.ExecutionGraphTest;
import io.ray.streaming.runtime.master.JobRuntimeContext;
import io.ray.streaming.runtime.master.graphmanager.GraphManager;
import io.ray.streaming.runtime.master.graphmanager.GraphManagerImpl;
import io.ray.streaming.runtime.master.scheduler.strategy.SlotAssignStrategy;
import io.ray.streaming.runtime.master.scheduler.strategy.impl.PipelineFirstStrategy;
import io.ray.streaming.runtime.master.resourcemanager.ResourceAssignmentView;
import io.ray.streaming.runtime.master.resourcemanager.strategy.ResourceAssignStrategy;
import io.ray.streaming.runtime.master.resourcemanager.strategy.impl.PipelineFirstStrategy;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import org.aeonbits.owner.ConfigFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.testng.Assert;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.Test;
@@ -33,63 +30,50 @@ public class PipelineFirstStrategyTest extends BaseUnitTest {
private Logger LOG = LoggerFactory.getLogger(PipelineFirstStrategyTest.class);
private SlotAssignStrategy strategy;
private List<Container> containers = new ArrayList<>();
private JobGraph jobGraph;
private ExecutionGraph executionGraph;
private int maxParallelism;
private ResourceAssignStrategy strategy;
@BeforeClass
public void setUp() {
strategy = new PipelineFirstStrategy();
Map<String, String> conf = new HashMap<>();
ResourceConfig resourceConfig = ConfigFactory.create(ResourceConfig.class, conf);
Resources resources = new Resources(resourceConfig);
Map<String, Double> containerResource = new HashMap<>();
containerResource.put(ResourceType.CPU.name(), 16.0);
containerResource.put(ResourceType.MEM.name(), 128.0);
for (int i = 0; i < 2; ++i) {
UniqueId uniqueId = UniqueId.randomId();
Container container = new Container(uniqueId, "1.1.1." + i, "localhost" + i);
container.setAvailableResource(containerResource);
Map<String, Double> resource = new HashMap<>();
resource.put(ResourceType.CPU.getValue(), 4.0);
resource.put(ResourceType.MEM.getValue(), 4.0);
Container container = new Container("1.1.1." + i, uniqueId, "localhost" + i, resource);
container.getAvailableResources().put(container.getName(), 500.0);
containers.add(container);
resources.getRegisterContainers().add(container);
}
strategy.setResources(resources);
}
//build ExecutionGraph
@AfterMethod
public void tearDown() {
reset();
}
private void reset() {
containers = new ArrayList<>();
strategy = null;
}
@Test
public void testStrategyName() {
Assert
.assertEquals(ResourceAssignStrategyType.PIPELINE_FIRST_STRATEGY.getName(), strategy.getName());
}
@Test
public void testAssignResource() {
strategy = new PipelineFirstStrategy();
Map<String, String> jobConf = new HashMap<>();
StreamingConfig streamingConfig = new StreamingConfig(jobConf);
GraphManager graphManager = new GraphManagerImpl(new JobRuntimeContext(streamingConfig));
jobGraph = ExecutionGraphTest.buildJobGraph();
executionGraph = ExecutionGraphTest.buildExecutionGraph(graphManager, jobGraph);
maxParallelism = executionGraph.getMaxParallelism();
JobGraph jobGraph = ExecutionGraphTest.buildJobGraph();
ExecutionGraph executionGraph = ExecutionGraphTest.buildExecutionGraph(graphManager, jobGraph);
ResourceAssignmentView assignmentView = strategy.assignResource(containers, executionGraph);
Assert.assertNotNull(assignmentView);
}
@Test
public int testSlotNumPerContainer() {
int slotNumPerContainer = strategy.getSlotNumPerContainer(containers, maxParallelism);
Assert.assertEquals(slotNumPerContainer,
(int) Math.ceil(Math.max(maxParallelism, containers.size()) * 1.0 / containers.size()));
return slotNumPerContainer;
}
@Test
public void testAllocateSlot() {
int slotNumPerContainer = testSlotNumPerContainer();
strategy.allocateSlot(containers, slotNumPerContainer);
for (Container container : containers) {
Assert.assertEquals(container.getSlots().size(), slotNumPerContainer);
}
}
@Test
public void testAssignSlot() {
Map<ContainerID, List<Slot>> allocatingMap = strategy.assignSlot(executionGraph);
for (Entry<ContainerID, List<Slot>> containerSlotEntry : allocatingMap.entrySet()) {
containerSlotEntry.getValue()
.forEach(slot -> Assert.assertNotEquals(slot.getExecutionVertexIds().size(), 0));
}
}
}
@@ -0,0 +1,82 @@
package io.ray.streaming.runtime.util;
import io.ray.api.id.UniqueId;
import io.ray.api.runtimecontext.NodeInfo;
import io.ray.streaming.runtime.core.resource.ResourceType;
import java.util.HashMap;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import org.powermock.api.mockito.PowerMockito;
/**
* Mockitools is a tool based on powermock and mokito to mock external service api
*/
public class Mockitools {
/**
* Mock GCS get node info api
*/
public static void mockGscApi() {
PowerMockito.mockStatic(RayUtils.class);
PowerMockito.when(RayUtils.getAliveNodeInfoMap())
.thenReturn(mockGetNodeInfoMap(mockGetAllNodeInfo()));
}
/**
* Mock get all node info from GCS
* @return
*/
public static List<NodeInfo> mockGetAllNodeInfo() {
List<NodeInfo> nodeInfos = new LinkedList<>();
for (int i = 1; i <= 5; i++) {
Map<String, Double> resources = new HashMap<>();
resources.put("MEM", 16.0);
switch (i) {
case 1:
resources.put(ResourceType.CPU.getValue(), 3.0);
break;
case 2:
case 3:
case 4:
resources.put(ResourceType.CPU.getValue(), 4.0);
break;
case 5:
resources.put(ResourceType.CPU.getValue(), 2.0);
break;
}
nodeInfos.add(mockNodeInfo(i, resources));
}
return nodeInfos;
}
/**
* Mock get node info map
* @param nodeInfos all node infos fetched from GCS
* @return node info map, key is node unique id, value is node info
*/
public static Map<UniqueId, NodeInfo> mockGetNodeInfoMap(List<NodeInfo> nodeInfos) {
return nodeInfos.stream().filter(nodeInfo -> nodeInfo.isAlive).collect(
Collectors.toMap(nodeInfo -> nodeInfo.nodeId, nodeInfo -> nodeInfo));
}
private static NodeInfo mockNodeInfo(int i, Map<String, Double> resources) {
return new NodeInfo(
createNodeId(i),
"localhost" + i,
"localhost" + i,
true,
resources);
}
private static UniqueId createNodeId(int id) {
byte[] nodeIdBytes = new byte[UniqueId.LENGTH];
for (int byteIndex = 0; byteIndex < UniqueId.LENGTH; ++byteIndex) {
nodeIdBytes[byteIndex] = String.valueOf(id).getBytes()[0];
}
return new UniqueId(nodeIdBytes);
}
}