Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<TaskName, List<SystemStreamPartition>> expectedTaskToSSPAssignments = ImmutableMap.of(new TaskName("task-0"), ImmutableList.of(testSystemStreamPartition1),
Expand All @@ -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<String, Map<String, String>> 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<String, LocationId> processorLocality = JobModelManager.getProcessorLocality(config, mockLocalityManager);

Mockito.verify(mockLocalityManager).readContainerLocality();
ImmutableMap<String, LocationId> 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<String, Map<String, String>> 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<String, LocationId> processorLocality = JobModelManager.getProcessorLocality(new MapConfig(), mockLocalityManager);
Map<String, LocationId> 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<String, LocationId> 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
Expand Down