From 548d18e75053ecf368bd695aeb1de77c839a0099 Mon Sep 17 00:00:00 2001 From: Cameron Lee Date: Mon, 19 Aug 2019 14:41:33 -0700 Subject: [PATCH] SAMZA-2304: Existing container locality mapping is incorrect when building job model --- .../samza/coordinator/JobModelManager.scala | 2 +- .../coordinator/TestJobModelManager.java | 37 +++++++++++++++---- 2 files changed, 31 insertions(+), 8 deletions(-) diff --git a/samza-core/src/main/scala/org/apache/samza/coordinator/JobModelManager.scala b/samza-core/src/main/scala/org/apache/samza/coordinator/JobModelManager.scala index 3414e8f35d..895ad00649 100644 --- a/samza-core/src/main/scala/org/apache/samza/coordinator/JobModelManager.scala +++ b/samza-core/src/main/scala/org/apache/samza/coordinator/JobModelManager.scala @@ -171,7 +171,7 @@ object JobModelManager extends Logging { val containerToLocationId: util.Map[String, LocationId] = new util.HashMap[String, LocationId]() val existingContainerLocality = localityManager.readContainerLocality() - for (containerId <- 0 to new JobConfig(config).getContainerCount) { + for (containerId <- 0 until new JobConfig(config).getContainerCount) { val localityMapping = existingContainerLocality.get(containerId.toString) // To handle the case when the container count is increased between two different runs of a samza-yarn job, // set the locality of newly added containers to any_host. diff --git a/samza-core/src/test/java/org/apache/samza/coordinator/TestJobModelManager.java b/samza-core/src/test/java/org/apache/samza/coordinator/TestJobModelManager.java index 0e53a5ce90..a2026eb13e 100644 --- a/samza-core/src/test/java/org/apache/samza/coordinator/TestJobModelManager.java +++ b/samza-core/src/test/java/org/apache/samza/coordinator/TestJobModelManager.java @@ -27,6 +27,7 @@ import com.google.common.collect.ImmutableSet; import org.apache.samza.Partition; import org.apache.samza.config.Config; +import org.apache.samza.config.JobConfig; import org.apache.samza.config.MapConfig; import org.apache.samza.container.LocalityManager; import org.apache.samza.container.TaskName; @@ -197,7 +198,7 @@ public void testGetGrouperMetadata() { Mockito.verify(mockLocalityManager).readContainerLocality(); Mockito.verify(mockTaskAssignmentManager).readTaskAssignment(); - Assert.assertEquals(ImmutableMap.of("0", new LocationId("abc-affinity"), "1", new LocationId("ANY_HOST")), grouperMetadata.getProcessorLocality()); + Assert.assertEquals(ImmutableMap.of("0", new LocationId("abc-affinity")), grouperMetadata.getProcessorLocality()); Assert.assertEquals(ImmutableMap.of(new TaskName("task-0"), new LocationId("abc-affinity")), grouperMetadata.getTaskLocality()); Map> expectedTaskToSSPAssignments = ImmutableMap.of(new TaskName("task-0"), ImmutableList.of(testSystemStreamPartition1), @@ -208,20 +209,42 @@ public void testGetGrouperMetadata() { } @Test - public void testGetProcessorLocality() { - // Mock the dependencies. + public void testGetProcessorLocalityAllEntriesExisting() { + Config config = new MapConfig(ImmutableMap.of(JobConfig.JOB_CONTAINER_COUNT, "2")); + + Map> localityMappings = new HashMap<>(); + localityMappings.put("0", ImmutableMap.of(SetContainerHostMapping.HOST_KEY, "0-affinity")); + localityMappings.put("1", ImmutableMap.of(SetContainerHostMapping.HOST_KEY, "1-affinity")); LocalityManager mockLocalityManager = mock(LocalityManager.class); + when(mockLocalityManager.readContainerLocality()).thenReturn(localityMappings); + + Map processorLocality = JobModelManager.getProcessorLocality(config, mockLocalityManager); + + Mockito.verify(mockLocalityManager).readContainerLocality(); + ImmutableMap expected = + ImmutableMap.of("0", new LocationId("0-affinity"), "1", new LocationId("1-affinity")); + Assert.assertEquals(expected, processorLocality); + } + + @Test + public void testGetProcessorLocalityNewContainer() { + Config config = new MapConfig(ImmutableMap.of(JobConfig.JOB_CONTAINER_COUNT, "2")); Map> localityMappings = new HashMap<>(); + // 2 containers, but only return 1 existing mapping localityMappings.put("0", ImmutableMap.of(SetContainerHostMapping.HOST_KEY, "abc-affinity")); - - // Mock the container locality assignment. + LocalityManager mockLocalityManager = mock(LocalityManager.class); when(mockLocalityManager.readContainerLocality()).thenReturn(localityMappings); - Map processorLocality = JobModelManager.getProcessorLocality(new MapConfig(), mockLocalityManager); + Map processorLocality = JobModelManager.getProcessorLocality(config, mockLocalityManager); Mockito.verify(mockLocalityManager).readContainerLocality(); - Assert.assertEquals(ImmutableMap.of("0", new LocationId("abc-affinity"), "1", new LocationId("ANY_HOST")), processorLocality); + ImmutableMap expected = ImmutableMap.of( + // found entry in existing locality + "0", new LocationId("abc-affinity"), + // no entry in existing locality + "1", new LocationId("ANY_HOST")); + Assert.assertEquals(expected, processorLocality); } @Test