From a9fa828948ea430bf9eff8338023f2a2d54e1395 Mon Sep 17 00:00:00 2001 From: Sanil15 Date: Thu, 12 Sep 2019 16:47:20 -0700 Subject: [PATCH 1/6] SAMZA-2319: Simplify Container Allocation logic --- .../AbstractContainerAllocator.java | 190 ++++++++++++++++-- .../clustermanager/ContainerAllocator.java | 74 ------- .../ContainerProcessManager.java | 22 +- .../HostAwareContainerAllocator.java | 177 ---------------- .../MockContainerAllocator.java | 5 +- .../MockHostAwareContainerAllocator.java | 7 +- .../TestContainerAllocator.java | 8 +- .../TestContainerProcessManager.java | 11 +- .../TestHostAwareContainerAllocator.java | 8 +- 9 files changed, 198 insertions(+), 304 deletions(-) delete mode 100644 samza-core/src/main/java/org/apache/samza/clustermanager/ContainerAllocator.java delete mode 100644 samza-core/src/main/java/org/apache/samza/clustermanager/HostAwareContainerAllocator.java diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java index bba4b520b9..91a0819a37 100644 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java +++ b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java @@ -35,20 +35,35 @@ /** * {@link AbstractContainerAllocator} makes requests for physical resources to the resource manager and also runs - * a processor on an allocated physical resource. Sub-classes should override the assignResourceRequests() - * method to assign resource requests according to some strategy. + * a processor on an allocated physical resource. * - * See {@link ContainerAllocator} and {@link HostAwareContainerAllocator} for two such strategies + * In case of host-affinity enabled, each request ({@link SamzaResourceRequest} encapsulates the identifier of the processor + * to be run and a "preferredHost". preferredHost is determined by the locality mappings in the coordinator stream. + * This thread periodically wakes up and makes the best-effort to assign a processor to the preferredHost. If a + * resource on the preferredHost is not returned by the cluster manager before the corresponding request expires, it + * assigns the processor to any other host that is allocated next. + * + * When host-affinity is not enabled, this periodically wakes up to assign a processor to *ANY* allocated resource. + * If there aren't enough resources, it waits by sleeping for {@code allocatorSleepIntervalMs} milliseconds. + * + * The resource expiry timeout is determined by CONTAINER_REQUEST_TIMEOUT and is configurable on a per-job basis. + * + * If there aren't enough resources, it waits by sleeping for allocatorSleepIntervalMs milliseconds. * * This class is not thread-safe. This class is used in the refactored code path as called by run-jc.sh */ -public abstract class AbstractContainerAllocator implements Runnable { +public class AbstractContainerAllocator implements Runnable { - private static final Logger log = LoggerFactory.getLogger(AbstractContainerAllocator.class); + private static final Logger LOG = LoggerFactory.getLogger(AbstractContainerAllocator.class); - /* State that controls the lifecycle of the allocator thread*/ + /* State that controls the lifecycle of the allocator thread */ private volatile boolean isRunning = true; + /** + * Flag for affine host requests + */ + private final boolean hostAffinityEnabled; + /** * Config and derived config objects */ @@ -83,22 +98,32 @@ public abstract class AbstractContainerAllocator implements Runnable { * ResourceRequestState indicates the state of all unfulfilled and allocated container requests */ protected final ResourceRequestState resourceRequestState; + /** + * Tracks the expiration of a request for resources. + */ + private final int requestTimeout; - public AbstractContainerAllocator(ClusterResourceManager containerProcessManager, - ResourceRequestState resourceRequestState, + private final Optional standbyContainerManager; + + public AbstractContainerAllocator(ClusterResourceManager clusterResourceManager, Config config, SamzaApplicationState state, - ClassLoader pluginClassLoader) { + ClassLoader pluginClassLoader, + boolean hostAffinityEnabled, + Optional standbyContainerManager) { ClusterManagerConfig clusterManagerConfig = new ClusterManagerConfig(config); - this.clusterResourceManager = containerProcessManager; + this.clusterResourceManager = clusterResourceManager; this.allocatorSleepIntervalMs = clusterManagerConfig.getAllocatorSleepTime(); - this.resourceRequestState = resourceRequestState; + this.resourceRequestState = new ResourceRequestState(hostAffinityEnabled, this.clusterResourceManager); this.containerMemoryMb = clusterManagerConfig.getContainerMemoryMb(); this.containerNumCpuCores = clusterManagerConfig.getNumCores(); this.taskConfig = new TaskConfig(config); this.state = state; this.config = config; this.pluginClassLoader = pluginClassLoader; + this.hostAffinityEnabled = hostAffinityEnabled; + this.standbyContainerManager = standbyContainerManager; + this.requestTimeout = clusterManagerConfig.getContainerRequestTimeout(); } /** @@ -121,18 +146,114 @@ public void run() { Thread.sleep(allocatorSleepIntervalMs); } catch (InterruptedException e) { - log.warn("Got InterruptedException in AllocatorThread.", e); + LOG.warn("Got InterruptedException in AllocatorThread.", e); Thread.currentThread().interrupt(); } catch (Exception e) { - log.error("Got unknown Exception in AllocatorThread.", e); + LOG.error("Got unknown Exception in AllocatorThread.", e); } } } /** * Assigns resources received from the cluster manager to processors. + * + * During the run() method, the thread sleeps for allocatorSleepIntervalMs ms. It then invokes assignResourceRequests, + * and tries to allocate any unsatisfied request that is still in the request queue {@link ResourceRequestState}) + * with allocated resources. + * When {@code hostAffinityEnabled} is disabled, all allocated resources are buffered in the list keyed by "ANY_HOST". + * When {@code hostAffinityEnabled} is enabled, all allocated resources are buffered in the list keyed by "preferredHost + * + * If the requested host is not available, the thread checks to see if the request has expired. If it has expired + * then two cases are handled seperately + * + * Case 1: host-affinity is disabled, cancels the current request and issues another ANY_HOST request + * Case 2: host-affinity is enabled, looks for allocated resouces on ANY_HOST and issues a container start if available, + * otherwise issues an ANY_HOST request + * + * In either of the scenarious if a {@code StandbyContainerManager} is present, the allocator transfers the request + * to it for checking StandByConstraints */ - protected abstract void assignResourceRequests(); + protected void assignResourceRequests() { + while (hasReadyPendingRequest()) { + SamzaResourceRequest request = peekReadyPendingRequest().get(); + String processorId = request.getProcessorId(); + String preferredHost = hostAffinityEnabled ? request.getPreferredHost() : ResourceRequestState.ANY_HOST; + Instant requestCreationTime = request.getRequestTimestamp(); + + LOG.info("Handling assignment request for Processor ID: {} on host: {}.", processorId, preferredHost); + if (hasAllocatedResource(preferredHost)) { + + // Found allocated container on preferredHost + LOG.info("Found an available container for Processor ID: {} on the host: {}", processorId, preferredHost); + + // Needs to be only updated when host affinity is enabled + if (hostAffinityEnabled) { + state.matchedResourceRequests.incrementAndGet(); + } + + // Try to launch processor on this preferredHost if it all standby constraints are met + if (this.standbyContainerManager.isPresent()) { + standbyContainerManager.get().checkStandbyConstraintsAndRunStreamProcessor(request, preferredHost, + peekAllocatedResource(preferredHost), this, resourceRequestState); + } else { + runStreamProcessor(request, preferredHost); + } + + } else { + + LOG.info("Did not find any allocated containers for running Processor ID: {} on the host: {}.", + processorId, preferredHost); + boolean expired = isRequestExpired(request); + + if (expired) { + updateExpiryMetrics(request); + if (hostAffinityEnabled) { + handleExpiredRequestWithHostAffinityEnabled(processorId, preferredHost, request); + } else { + handleExpiredRequestWithHostAffinityDisabled(processorId, request); + } + } else { + LOG.info("Request for Processor ID: {} on host: {} has not expired yet." + + "Request creation time: {}. Current Time: {}. Request timeout: {} ms", processorId, preferredHost, + requestCreationTime, System.currentTimeMillis(), requestTimeout); + break; + } + } + } + } + + /** + * Handles an expired resource request when {@code hostAffinityEnabled} is true, in this case since the + * preferred host, we try to see if a surplus ANY_HOST is available in the request queue. + */ + private void handleExpiredRequestWithHostAffinityEnabled(String processorId, String preferredHost, + SamzaResourceRequest request) { + boolean resourceAvailableOnAnyHost = hasAllocatedResource(ResourceRequestState.ANY_HOST); + if (standbyContainerManager.isPresent()) { + standbyContainerManager.get() + .handleExpiredResourceRequest(processorId, request, + Optional.ofNullable(peekAllocatedResource(ResourceRequestState.ANY_HOST)), this, resourceRequestState); + } else if (resourceAvailableOnAnyHost) { + LOG.info("Request for Processor ID: {} on host: {} has expired. Running on ANY_HOST", processorId, preferredHost); + runStreamProcessor(request, ResourceRequestState.ANY_HOST); + } else { + LOG.info("Request for Processor ID: {} on host: {} has expired. Requesting additional resources on ANY_HOST.", + processorId, preferredHost); + resourceRequestState.cancelResourceRequest(request); + requestResource(processorId, ResourceRequestState.ANY_HOST); + } + } + + /** + * Handles an expired resource request when {@code hostAffinityEnabled} is false, in this case since there are no + * ANY_HOST available in the request queue we cancel existing request & reissue a new one + */ + private void handleExpiredRequestWithHostAffinityDisabled(String processorId, SamzaResourceRequest request) { + LOG.info("Request for Processor ID: {} on ANY_HOST has expired. Requesting additional resources on ANY_HOST.", + processorId); + resourceRequestState.cancelResourceRequest(request); + requestResource(processorId, ResourceRequestState.ANY_HOST); + } /** * Updates the request state and runs a processor on the specified host. Assumes a resource @@ -156,7 +277,7 @@ protected void runStreamProcessor(SamzaResourceRequest request, String preferred String processorId = request.getProcessorId(); // Run processor on resource - log.info("Found Container ID: {} for Processor ID: {} on host: {} for request creation time: {}.", + LOG.info("Found Container ID: {} for Processor ID: {} on host: {} for request creation time: {}.", resource.getContainerId(), processorId, preferredHost, request.getRequestTimestamp()); // Update processor state as "pending" and then issue a request to launch it. It's important to perform the state-update @@ -178,7 +299,19 @@ protected void runStreamProcessor(SamzaResourceRequest request, String preferred * - when host-affinity is enabled and job is run for the first time * - when the number of containers has been increased. */ - public abstract void requestResources(Map processorToHostMapping); + public void requestResources(Map processorToHostMapping) { + for (Map.Entry entry : processorToHostMapping.entrySet()) { + String processorId = entry.getKey(); + String preferredHost = entry.getValue(); + if (!hostAffinityEnabled) { + preferredHost = ResourceRequestState.ANY_HOST; + } else if (preferredHost == null) { + LOG.info("No preferred host mapping found for Processor ID: {}. Requesting resource on ANY_HOST", processorId); + preferredHost = ResourceRequestState.ANY_HOST; + } + requestResource(processorId, preferredHost); + } + } /** * Checks if this allocator has a pending resource request with a request timestamp equal to or earlier than the current @@ -309,4 +442,29 @@ public final void releaseResource(String containerId) { public void stop() { isRunning = false; } + + + /** + * Checks if a request has expired. + * @param request the request to check + * @return true if request has expired + */ + private boolean isRequestExpired(SamzaResourceRequest request) { + long currTime = Instant.now().toEpochMilli(); + boolean requestExpired = currTime - request.getRequestTimestamp().toEpochMilli() > requestTimeout; + if (requestExpired) { + LOG.info("Request for Processor ID: {} on host: {} with creation time: {} has expired at current time: {} after timeout: {} ms.", + request.getProcessorId(), request.getPreferredHost(), request.getRequestTimestamp(), currTime, requestTimeout); + } + return requestExpired; + } + + private void updateExpiryMetrics(SamzaResourceRequest request) { + String preferredHost = request.getPreferredHost(); + if (ResourceRequestState.ANY_HOST.equals(preferredHost)) { + state.expiredAnyHostRequests.incrementAndGet(); + } else { + state.expiredPreferredHostRequests.incrementAndGet(); + } + } } diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerAllocator.java deleted file mode 100644 index 9078f25a1f..0000000000 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerAllocator.java +++ /dev/null @@ -1,74 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.samza.clustermanager; - -import java.util.Map; -import org.apache.samza.config.Config; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * This is the default allocator that will be used by ContainerProcessManager. - * - * When host-affinity is not enabled, this periodically wakes up to assign a processor to *ANY* allocated resource. - * If there aren't enough resources, it waits by sleeping for {@code allocatorSleepIntervalMs} milliseconds. - * - * This class is instantiated by the ContainerProcessManager (which in turn is created by the JC from run-jc.sh), - * when host-affinity is off. Otherwise, the HostAwareContainerAllocator is instantiated. - */ -public class ContainerAllocator extends AbstractContainerAllocator { - private static final Logger log = LoggerFactory.getLogger(ContainerAllocator.class); - - public ContainerAllocator(ClusterResourceManager manager, - Config config, - SamzaApplicationState state, - ClassLoader pluginClassloader) { - super(manager, new ResourceRequestState(false, manager), config, state, pluginClassloader); - } - - /** - * During the run() method, the thread sleeps for allocatorSleepIntervalMs ms. It then invokes assignResourceRequests, - * and tries to allocate any unsatisfied request that is still in the request queue {@link ResourceRequestState}) - * with allocated resources, if any. - * - * Since host-affinity is not enabled, all allocated resources are buffered in the list keyed by "ANY_HOST". - * */ - @Override - public void assignResourceRequests() { - while (hasReadyPendingRequest() && hasAllocatedResource(ResourceRequestState.ANY_HOST)) { - peekReadyPendingRequest().ifPresent(request -> runStreamProcessor(request, ResourceRequestState.ANY_HOST)); - } - } - - /** - * Since host-affinity is not enabled, the processor id to host mappings will be ignored and all resources will be - * matched to any available host. - * - * @param processorToHostMapping A Map of [processorId, hostName] ID of the processor to run on the resource. - * The hostName will be ignored and each processor will be matched to any available host. - */ - @Override - public void requestResources(Map processorToHostMapping) { - for (Map.Entry entry : processorToHostMapping.entrySet()) { - String processorId = entry.getKey(); - requestResource(processorId, ResourceRequestState.ANY_HOST); - } - } - -} diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerProcessManager.java b/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerProcessManager.java index 8e12f825a4..a12fc3494c 100644 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerProcessManager.java +++ b/samza-core/src/main/java/org/apache/samza/clustermanager/ContainerProcessManager.java @@ -60,7 +60,7 @@ * - Identifying the cause of container failure and re-requesting containers from the cluster manager by adding request to the * internal requestQueue in {@link ResourceRequestState} * - The allocator thread that assigns the allocated containers to pending requests - * (See {@link org.apache.samza.clustermanager.ContainerAllocator} or {@link org.apache.samza.clustermanager.HostAwareContainerAllocator}) + * (See {@link org.apache.samza.clustermanager.AbstractContainerAllocator} or {@link org.apache.samza.clustermanager.AbstractContainerAllocator}) * */ public class ContainerProcessManager implements ClusterResourceManager.Callback { @@ -167,9 +167,7 @@ public ContainerProcessManager(Config config, SamzaApplicationState state, Metri this.standbyContainerManager = Optional.empty(); } - this.containerAllocator = - buildContainerAllocator(this.hostAffinityEnabled, this.clusterResourceManager, this.clusterManagerConfig, - config, this.standbyContainerManager, state, classLoader); + this.containerAllocator = new AbstractContainerAllocator(this.clusterResourceManager, config, state, classLoader, hostAffinityEnabled, this.standbyContainerManager); this.allocatorThread = new Thread(this.containerAllocator, "Container Allocator Thread"); LOG.info("Finished container process manager initialization."); } @@ -191,24 +189,12 @@ public ContainerProcessManager(Config config, SamzaApplicationState state, Metri this.standbyContainerManager = Optional.empty(); this.diagnosticsManager = Option.empty(); this.containerAllocator = allocator.orElseGet( - () -> buildContainerAllocator(this.hostAffinityEnabled, this.clusterResourceManager, this.clusterManagerConfig, - clusterManagerConfig, this.standbyContainerManager, state, classLoader)); - + () -> new AbstractContainerAllocator(this.clusterResourceManager, clusterManagerConfig, state, classLoader, + hostAffinityEnabled, this.standbyContainerManager)); this.allocatorThread = new Thread(this.containerAllocator, "Container Allocator Thread"); LOG.info("Finished container process manager initialization"); } - private static AbstractContainerAllocator buildContainerAllocator(boolean hostAffinityEnabled, - ClusterResourceManager clusterResourceManager, ClusterManagerConfig clusterManagerConfig, Config config, - Optional standbyContainerManager, SamzaApplicationState state, ClassLoader classLoader) { - if (hostAffinityEnabled) { - return new HostAwareContainerAllocator(clusterResourceManager, clusterManagerConfig.getContainerRequestTimeout(), - config, standbyContainerManager, state, classLoader); - } else { - return new ContainerAllocator(clusterResourceManager, config, state, classLoader); - } - } - public boolean shouldShutdown() { LOG.debug("ContainerProcessManager state: Completed containers: {}, Configured containers: {}, Are there too many failed containers: {}, Is allocator thread alive: {}", state.completedProcessors.get(), state.processorCount, jobFailureCriteriaMet ? "yes" : "no", allocatorThread.isAlive() ? "yes" : "no"); diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/HostAwareContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/HostAwareContainerAllocator.java deleted file mode 100644 index 177c473920..0000000000 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/HostAwareContainerAllocator.java +++ /dev/null @@ -1,177 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.samza.clustermanager; - -import java.time.Instant; -import java.util.Map; -import java.util.Optional; -import org.apache.samza.config.Config; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - - -/** - * This is the allocator thread that will be used by ContainerProcessManager when host-affinity is enabled for a job. It is similar - * to {@link ContainerAllocator}, except that it considers locality for allocation. - * - * In case of host-affinity, each request ({@link SamzaResourceRequest} encapsulates the identifier of the processor - * to be run and a "preferredHost". preferredHost is determined by the locality mappings in the coordinator stream. - * This thread periodically wakes up and makes the best-effort to assign a processor to the preferredHost. If a - * resource on the preferredHost is not returned by the cluster manager before the corresponding request expires, it - * assigns the processor to any other host that is allocated next. - * - * The resource expiry timeout is determined by CONTAINER_REQUEST_TIMEOUT and is configurable on a per-job basis. - * - * If there aren't enough resources, it waits by sleeping for allocatorSleepIntervalMs milliseconds. - */ -//This class is used in the refactored code path as called by run-jc.sh -public class HostAwareContainerAllocator extends AbstractContainerAllocator { - private static final Logger log = LoggerFactory.getLogger(HostAwareContainerAllocator.class); - /** - * Tracks the expiration of a request for resources. - */ - private final int requestTimeout; - private final Optional standbyContainerManager; - - public HostAwareContainerAllocator(ClusterResourceManager manager, - int timeout, - Config config, - Optional standbyContainerManager, - SamzaApplicationState state, - ClassLoader pluginClassLoader) { - super(manager, new ResourceRequestState(true, manager), config, state, pluginClassLoader); - this.requestTimeout = timeout; - this.standbyContainerManager = standbyContainerManager; - } - - /** - * Since host-affinity is enabled, all allocated resources are buffered in the list keyed by "preferredHost". - * - * If the requested host is not available, the thread checks to see if the request has expired. - * If it has expired, it runs the processor on one of the available resources from the - * allocatedContainers buffer keyed by "ANY_HOST". - */ - @Override - public void assignResourceRequests() { - for (Optional requestOptional = peekReadyPendingRequest(); - requestOptional.isPresent(); - requestOptional = peekReadyPendingRequest()) { - SamzaResourceRequest request = requestOptional.get(); - String processorId = request.getProcessorId(); - String preferredHost = request.getPreferredHost(); - Instant requestCreationTime = request.getRequestTimestamp(); - log.info("Handling assignment request for Processor ID: {} on preferred host: {}.", processorId, preferredHost); - - if (hasAllocatedResource(preferredHost)) { - // Found allocated container on preferredHost - log.info("Found an available container for Processor ID: {} on the preferred host: {}", processorId, preferredHost); - // Try to launch processor on this preferredHost if it all standby constraints are met - checkStandbyConstraintsAndRunStreamProcessor(request, preferredHost, peekAllocatedResource(preferredHost)); - state.matchedResourceRequests.incrementAndGet(); - } else { - log.info("Did not find any allocated containers for running Processor ID: {} on the preferred host: {}.", processorId, preferredHost); - - boolean expired = isRequestExpired(request); - boolean resourceAvailableOnAnyHost = hasAllocatedResource(ResourceRequestState.ANY_HOST); - - if (expired) { - updateExpiryMetrics(request); - - if (standbyContainerManager.isPresent()) { - standbyContainerManager.get().handleExpiredResourceRequest(processorId, request, - Optional.ofNullable(peekAllocatedResource(ResourceRequestState.ANY_HOST)), this, resourceRequestState); - - } else if (resourceAvailableOnAnyHost) { - log.info("Request for Processor ID: {} on host: {} has expired. Running on ANY_HOST", processorId, preferredHost); - runStreamProcessor(request, ResourceRequestState.ANY_HOST); - - } else { - log.info("Request for Processor ID: {} on host: {} has expired. Requesting additional resources on ANY_HOST.", processorId, preferredHost); - resourceRequestState.cancelResourceRequest(request); - requestResource(processorId, ResourceRequestState.ANY_HOST); - } - - } else { - log.info("Request for Processor ID: {} on host: {} has not expired yet." + - "Request creation time: {}. Current Time: {}. Request timeout: {} ms", - processorId, preferredHost, requestCreationTime, System.currentTimeMillis(), requestTimeout); - break; - } - } - } - } - - - /** - * Since host-affinity is enabled, containers for all processors will be requested on their preferred host. If the job is - * run for the first time, it will get matched to any available host. - * - * @param processorToHostMapping A Map of [processorId, hostName] where processorId is the ID of the Samza processor - * to run on the resource. hostName is the host on which the resource must be allocated. - * The hostName value is null when host-affinity is enabled and job is run for the - * first time, or when the number of containers has been increased. - */ - @Override - public void requestResources(Map processorToHostMapping) { - for (Map.Entry entry : processorToHostMapping.entrySet()) { - String processorId = entry.getKey(); - String preferredHost = entry.getValue(); - if (preferredHost == null) { - log.info("No preferred host mapping found for Processor ID: {}. Requesting resource on ANY_HOST", processorId); - preferredHost = ResourceRequestState.ANY_HOST; - } - requestResource(processorId, preferredHost); - } - } - - /** - * Checks if a request has expired. - * @param request the request to check - * @return true if request has expired - */ - private boolean isRequestExpired(SamzaResourceRequest request) { - long currTime = Instant.now().toEpochMilli(); - boolean requestExpired = currTime - request.getRequestTimestamp().toEpochMilli() > requestTimeout; - if (requestExpired) { - log.info("Request for Processor ID: {} on host: {} with creation time: {} has expired at current time: {} after timeout: {} ms.", - request.getProcessorId(), request.getPreferredHost(), request.getRequestTimestamp(), currTime, requestTimeout); - } - return requestExpired; - } - - private void updateExpiryMetrics(SamzaResourceRequest request) { - String preferredHost = request.getPreferredHost(); - if (ResourceRequestState.ANY_HOST.equals(preferredHost)) { - state.expiredAnyHostRequests.incrementAndGet(); - } else { - state.expiredPreferredHostRequests.incrementAndGet(); - } - } - - private void checkStandbyConstraintsAndRunStreamProcessor(SamzaResourceRequest request, String preferredHost, SamzaResource samzaResource) { - // If standby tasks are not enabled run streamprocessor on the given host - if (!this.standbyContainerManager.isPresent()) { - runStreamProcessor(request, preferredHost); - return; - } - - this.standbyContainerManager.get().checkStandbyConstraintsAndRunStreamProcessor(request, preferredHost, - samzaResource, this, resourceRequestState); - } -} \ No newline at end of file diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java index e294cfc5ad..dae8c05583 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java @@ -18,6 +18,7 @@ */ package org.apache.samza.clustermanager; +import java.util.Optional; import org.apache.samza.config.Config; import java.lang.reflect.Field; @@ -26,14 +27,14 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; -public class MockContainerAllocator extends ContainerAllocator { +public class MockContainerAllocator extends AbstractContainerAllocator { public int requestedContainers = 0; private Semaphore semaphore = new Semaphore(0); public MockContainerAllocator(ClusterResourceManager manager, Config config, SamzaApplicationState state) { - super(manager, config, state, MockContainerAllocator.class.getClassLoader()); + super(manager, config, state, MockContainerAllocator.class.getClassLoader(), false, Optional.empty()); } /** diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java index 0c44edf2ee..9e4bf8ce52 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java @@ -26,14 +26,13 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; -public class MockHostAwareContainerAllocator extends HostAwareContainerAllocator { - private static final int ALLOCATOR_TIMEOUT_MS = 10000; +public class MockHostAwareContainerAllocator extends AbstractContainerAllocator { private Semaphore semaphore = new Semaphore(0); public MockHostAwareContainerAllocator(ClusterResourceManager manager, Config config, SamzaApplicationState state) { - super(manager, ALLOCATOR_TIMEOUT_MS, config, Optional.empty(), state, - MockHostAwareContainerAllocator.class.getClassLoader()); + super(manager, config, state, + MockHostAwareContainerAllocator.class.getClassLoader(), true, Optional.empty()); } /** diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java index 0f90f92b80..9083ad356c 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java @@ -22,6 +22,7 @@ import java.lang.reflect.Field; import java.util.HashMap; import java.util.Map; +import java.util.Optional; import org.apache.samza.config.Config; import org.apache.samza.config.MapConfig; import org.apache.samza.coordinator.JobModelManager; @@ -47,15 +48,16 @@ public class TestContainerAllocator { private final SamzaApplicationState state = new SamzaApplicationState(jobModelManager); private final MockClusterResourceManager manager = new MockClusterResourceManager(callback, state); - private ContainerAllocator containerAllocator; + private AbstractContainerAllocator containerAllocator; private MockContainerRequestState requestState; private Thread allocatorThread; @Before public void setup() throws Exception { - containerAllocator = new ContainerAllocator(manager, config, state, getClass().getClassLoader()); + containerAllocator = new AbstractContainerAllocator(manager, config, state, getClass().getClassLoader(), false, Optional + .empty()); requestState = new MockContainerRequestState(manager, false); - Field requestStateField = containerAllocator.getClass().getSuperclass().getDeclaredField("resourceRequestState"); + Field requestStateField = containerAllocator.getClass().getDeclaredField("resourceRequestState"); requestStateField.setAccessible(true); requestStateField.set(containerAllocator, requestState); allocatorThread = new Thread(containerAllocator); diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java index 997adb1c5d..f3f8ffeeb8 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java @@ -71,7 +71,7 @@ public class TestContainerProcessManager { put("cluster-manager.container.retry.count", "1"); put("cluster-manager.container.retry.window.ms", "1999999999"); put("cluster-manager.allocator.sleep.ms", "1"); - put("cluster-manager.container.request.timeout.ms", "2"); + put("cluster-manager.container.request.timeout.ms", "5000"); put("cluster-manager.container.memory.mb", "512"); put("yarn.package.path", "/foo"); put("task.inputs", "test-system.test-stream"); @@ -147,7 +147,7 @@ public void testContainerProcessManager() throws Exception { AbstractContainerAllocator allocator = (AbstractContainerAllocator) getPrivateFieldFromCpm("containerAllocator", cpm).get(cpm); - assertEquals(ContainerAllocator.class, allocator.getClass()); + assertEquals(AbstractContainerAllocator.class, allocator.getClass()); // Asserts that samza exposed container configs is honored by allocator thread assertEquals(500, allocator.containerMemoryMb); assertEquals(5, allocator.containerNumCpuCores); @@ -171,7 +171,7 @@ public void testContainerProcessManager() throws Exception { allocator = (AbstractContainerAllocator) getPrivateFieldFromCpm("containerAllocator", cpm).get(cpm); - assertEquals(HostAwareContainerAllocator.class, allocator.getClass()); + assertEquals(AbstractContainerAllocator.class, allocator.getClass()); // Asserts that samza exposed container configs is honored by allocator thread assertEquals(500, allocator.containerMemoryMb); assertEquals(5, allocator.containerNumCpuCores); @@ -630,6 +630,7 @@ public void testAllBufferedResourcesAreUtilized() throws Exception { config.putAll(getConfigWithHostAffinity()); config.put("job.container.count", "2"); config.put("cluster-manager.container.retry.count", "2"); + config.put("cluster-manager.container.request.timeout.ms", "10000"); Config cfg = new MapConfig(config); // 1. Request two containers on hosts - host1 and host2 SamzaApplicationState state = new SamzaApplicationState(getJobModelManagerWithHostAffinity(ImmutableMap.of("0", "host1", @@ -761,9 +762,7 @@ public void testDuplicateNotificationsDoNotAffectJobHealth() throws Exception { @Test public void testNewContainerRequestedOnFailureWithKnownCode() throws Exception { Config conf = getConfig(); - - Map config = new HashMap<>(); - config.putAll(getConfig()); + Map config = new HashMap<>(getConfig()); SamzaApplicationState state = new SamzaApplicationState(getJobModelManagerWithoutHostAffinity(1)); MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java index 40d292c017..29f173e259 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java @@ -68,7 +68,7 @@ private JobModelManager initializeJobModelManager(Config config, int containerCo new MockHttpServer("/", 7777, null, new ServletHolder(DefaultServlet.class))); } - private HostAwareContainerAllocator containerAllocator; + private AbstractContainerAllocator containerAllocator; private final int timeoutMillis = 1000; private MockContainerRequestState requestState; private Thread allocatorThread; @@ -76,10 +76,10 @@ private JobModelManager initializeJobModelManager(Config config, int containerCo @Before public void setup() throws Exception { containerAllocator = - new HostAwareContainerAllocator(clusterResourceManager, timeoutMillis, config, Optional.empty(), state, - getClass().getClassLoader()); + new AbstractContainerAllocator(clusterResourceManager, config, state, + getClass().getClassLoader(), true, Optional.empty()); requestState = new MockContainerRequestState(clusterResourceManager, true); - Field requestStateField = containerAllocator.getClass().getSuperclass().getDeclaredField("resourceRequestState"); + Field requestStateField = containerAllocator.getClass().getDeclaredField("resourceRequestState"); requestStateField.setAccessible(true); requestStateField.set(containerAllocator, requestState); allocatorThread = new Thread(containerAllocator); From 48a6215c554a15dbb84a8e21ad82de8aa3c022ab Mon Sep 17 00:00:00 2001 From: Sanil15 Date: Fri, 13 Sep 2019 15:41:11 -0700 Subject: [PATCH 2/6] Removing expiry check from ContainerAllocator for host affinity off case --- .../AbstractContainerAllocator.java | 15 +------------- ...ckContainerAllocatorWithHostAffinity.java} | 6 +++--- ...ontainerAllocatorWithoutHostAffinity.java} | 6 +++--- ...stContainerAllocatorWithHostAffinity.java} | 2 +- ...ontainerAllocatorWithoutHostAffinity.java} | 2 +- .../TestContainerProcessManager.java | 20 +++++++++---------- 6 files changed, 19 insertions(+), 32 deletions(-) rename samza-core/src/test/java/org/apache/samza/clustermanager/{MockHostAwareContainerAllocator.java => MockContainerAllocatorWithHostAffinity.java} (90%) rename samza-core/src/test/java/org/apache/samza/clustermanager/{MockContainerAllocator.java => MockContainerAllocatorWithoutHostAffinity.java} (89%) rename samza-core/src/test/java/org/apache/samza/clustermanager/{TestHostAwareContainerAllocator.java => TestContainerAllocatorWithHostAffinity.java} (99%) rename samza-core/src/test/java/org/apache/samza/clustermanager/{TestContainerAllocator.java => TestContainerAllocatorWithoutHostAffinity.java} (99%) diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java index 91a0819a37..95077500a5 100644 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java +++ b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java @@ -209,11 +209,9 @@ protected void assignResourceRequests() { updateExpiryMetrics(request); if (hostAffinityEnabled) { handleExpiredRequestWithHostAffinityEnabled(processorId, preferredHost, request); - } else { - handleExpiredRequestWithHostAffinityDisabled(processorId, request); } } else { - LOG.info("Request for Processor ID: {} on host: {} has not expired yet." + LOG.info("Request for Processor ID: {} on preferred host {} has not expired yet." + "Request creation time: {}. Current Time: {}. Request timeout: {} ms", processorId, preferredHost, requestCreationTime, System.currentTimeMillis(), requestTimeout); break; @@ -244,17 +242,6 @@ private void handleExpiredRequestWithHostAffinityEnabled(String processorId, Str } } - /** - * Handles an expired resource request when {@code hostAffinityEnabled} is false, in this case since there are no - * ANY_HOST available in the request queue we cancel existing request & reissue a new one - */ - private void handleExpiredRequestWithHostAffinityDisabled(String processorId, SamzaResourceRequest request) { - LOG.info("Request for Processor ID: {} on ANY_HOST has expired. Requesting additional resources on ANY_HOST.", - processorId); - resourceRequestState.cancelResourceRequest(request); - requestResource(processorId, ResourceRequestState.ANY_HOST); - } - /** * Updates the request state and runs a processor on the specified host. Assumes a resource * is available on the preferred host, so the caller must verify that before invoking this method. diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithHostAffinity.java similarity index 90% rename from samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java rename to samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithHostAffinity.java index 9e4bf8ce52..f7f62aab5f 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/MockHostAwareContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithHostAffinity.java @@ -26,13 +26,13 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; -public class MockHostAwareContainerAllocator extends AbstractContainerAllocator { +public class MockContainerAllocatorWithHostAffinity extends AbstractContainerAllocator { private Semaphore semaphore = new Semaphore(0); - public MockHostAwareContainerAllocator(ClusterResourceManager manager, + public MockContainerAllocatorWithHostAffinity(ClusterResourceManager manager, Config config, SamzaApplicationState state) { super(manager, config, state, - MockHostAwareContainerAllocator.class.getClassLoader(), true, Optional.empty()); + MockContainerAllocatorWithHostAffinity.class.getClassLoader(), true, Optional.empty()); } /** diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithoutHostAffinity.java similarity index 89% rename from samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java rename to samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithoutHostAffinity.java index dae8c05583..5ea4471f03 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithoutHostAffinity.java @@ -27,14 +27,14 @@ import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; -public class MockContainerAllocator extends AbstractContainerAllocator { +public class MockContainerAllocatorWithoutHostAffinity extends AbstractContainerAllocator { public int requestedContainers = 0; private Semaphore semaphore = new Semaphore(0); - public MockContainerAllocator(ClusterResourceManager manager, + public MockContainerAllocatorWithoutHostAffinity(ClusterResourceManager manager, Config config, SamzaApplicationState state) { - super(manager, config, state, MockContainerAllocator.class.getClassLoader(), false, Optional.empty()); + super(manager, config, state, MockContainerAllocatorWithoutHostAffinity.class.getClassLoader(), false, Optional.empty()); } /** diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java similarity index 99% rename from samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java rename to samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java index 29f173e259..7ccc27936f 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestHostAwareContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java @@ -47,7 +47,7 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; -public class TestHostAwareContainerAllocator { +public class TestContainerAllocatorWithHostAffinity { private final Config config = getConfig(); private final JobModelManager jobModelManager = initializeJobModelManager(config, 1); diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java similarity index 99% rename from samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java rename to samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java index 9083ad356c..978055a657 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocator.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java @@ -39,7 +39,7 @@ import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; -public class TestContainerAllocator { +public class TestContainerAllocatorWithoutHostAffinity { private final MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); private final Config config = getConfig(); private final JobModelManager jobModelManager = JobModelManagerTestUtil.getJobModelManager(config, 1, diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java index f3f8ffeeb8..7996e59944 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java @@ -188,7 +188,7 @@ public void testOnInit() throws Exception { ContainerProcessManager cpm = buildContainerProcessManager(clusterManagerConfig, state, clusterResourceManager, Optional.empty()); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, conf, state); @@ -249,7 +249,7 @@ public void testCpmShouldStopWhenContainersFinish() throws Exception { MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); ClusterManagerConfig clusterManagerConfig = spy(new ClusterManagerConfig(conf)); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, conf, state); @@ -293,7 +293,7 @@ public void testNewContainerRequestedOnFailureWithUnknownCode() throws Exception MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); ClusterManagerConfig clusterManagerConfig = spy(new ClusterManagerConfig(conf)); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, conf, state); @@ -386,7 +386,7 @@ private void testContainerRequestedRetriesExceedingWindowOnFailureWithUnknownCod MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, clusterManagerConfig, state); @@ -466,7 +466,7 @@ private void testContainerRequestedRetriesNotExceedingWindowOnFailureWithUnknown MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, clusterManagerConfig, state); @@ -568,7 +568,7 @@ public void testInvalidNotificationsAreIgnored() throws Exception { MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); ClusterManagerConfig clusterManagerConfig = spy(new ClusterManagerConfig(conf)); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, conf, state); @@ -606,7 +606,7 @@ public void testRerequestOnAnyHostIfContainerStartFails() throws Exception { MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, new MapConfig(config), state); @@ -638,7 +638,7 @@ public void testAllBufferedResourcesAreUtilized() throws Exception { MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); - MockHostAwareContainerAllocator allocator = new MockHostAwareContainerAllocator( + MockContainerAllocatorWithHostAffinity allocator = new MockContainerAllocatorWithHostAffinity( clusterResourceManager, cfg, state); @@ -699,7 +699,7 @@ public void testDuplicateNotificationsDoNotAffectJobHealth() throws Exception { MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); ClusterManagerConfig clusterManagerConfig = spy(new ClusterManagerConfig(new MapConfig(conf))); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, conf, state); @@ -768,7 +768,7 @@ public void testNewContainerRequestedOnFailureWithKnownCode() throws Exception { MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); ClusterManagerConfig clusterManagerConfig = spy(new ClusterManagerConfig(new MapConfig(config))); - MockContainerAllocator allocator = new MockContainerAllocator( + MockContainerAllocatorWithoutHostAffinity allocator = new MockContainerAllocatorWithoutHostAffinity( clusterResourceManager, conf, state); From f8a6b185271e5321e8f02ccd7dddb0ebf3a00b38 Mon Sep 17 00:00:00 2001 From: Sanil15 Date: Tue, 24 Sep 2019 10:49:34 -0700 Subject: [PATCH 3/6] Addressing Review Adding more tests --- .../AbstractContainerAllocator.java | 89 ++++++---- ...estContainerAllocatorWithHostAffinity.java | 161 +++++++++++++++++- ...ContainerAllocatorWithoutHostAffinity.java | 61 ++++++- .../TestContainerProcessManager.java | 6 +- .../diagnostics/TestDiagnosticsManager.java | 3 +- 5 files changed, 280 insertions(+), 40 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java index 95077500a5..c530d05af2 100644 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java +++ b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java @@ -18,6 +18,7 @@ */ package org.apache.samza.clustermanager; +import com.google.common.annotations.VisibleForTesting; import java.time.Duration; import java.time.Instant; import java.util.Map; @@ -37,18 +38,34 @@ * {@link AbstractContainerAllocator} makes requests for physical resources to the resource manager and also runs * a processor on an allocated physical resource. * - * In case of host-affinity enabled, each request ({@link SamzaResourceRequest} encapsulates the identifier of the processor - * to be run and a "preferredHost". preferredHost is determined by the locality mappings in the coordinator stream. - * This thread periodically wakes up and makes the best-effort to assign a processor to the preferredHost. If a - * resource on the preferredHost is not returned by the cluster manager before the corresponding request expires, it - * assigns the processor to any other host that is allocated next. - * - * When host-affinity is not enabled, this periodically wakes up to assign a processor to *ANY* allocated resource. - * If there aren't enough resources, it waits by sleeping for {@code allocatorSleepIntervalMs} milliseconds. - * - * The resource expiry timeout is determined by CONTAINER_REQUEST_TIMEOUT and is configurable on a per-job basis. - * - * If there aren't enough resources, it waits by sleeping for allocatorSleepIntervalMs milliseconds. + *
    + *
  • + * In case of host-affinity enabled, each request ({@link SamzaResourceRequest} contains a processorId which + * identifies the processor the request is for and a "preferredHost" which is determined by the locality mappings + * in the coordinator stream + *
  • + *
  • + * This thread periodically periodically matches outstanding resource requests with allocated resources. + * Its period is controlled using the {@code allocatorSleepIntervalMs} parameter + *
  • + *
  • + * In case of host-affinity is enabled, the resource-request's preferredHost param is set to the host the processor + * was last seen on + *
  • + *
  • + * In case of host-affinity us disabled, the resource-request's preferredHost param is set to {@link ResourceRequestState#ANY_HOST} + *
  • + *
  • + * When host-affinity is enabled and a preferred resource has not been obtained after {@code requestExpiryTimeout} + * milliseconds of the request being made, the resource is declared expired. The expired request are handled by + * allocating them to *ANY* allocated resource if available. If no surplus resources are available the current preferred + * resource-request is cancelled and resource-request for ANY_HOST is issued + *
  • + *
  • + * When host-affinity is not enabled, this periodically wakes up to assign a processor to *ANY* allocated resource. + * If there aren't enough resources, it waits by sleeping for {@code allocatorSleepIntervalMs} milliseconds. + *
  • + *
* * This class is not thread-safe. This class is used in the refactored code path as called by run-jc.sh */ @@ -101,7 +118,7 @@ public class AbstractContainerAllocator implements Runnable { /** * Tracks the expiration of a request for resources. */ - private final int requestTimeout; + private final int requestExpiryTimeout; private final Optional standbyContainerManager; @@ -123,7 +140,7 @@ public AbstractContainerAllocator(ClusterResourceManager clusterResourceManager, this.pluginClassLoader = pluginClassLoader; this.hostAffinityEnabled = hostAffinityEnabled; this.standbyContainerManager = standbyContainerManager; - this.requestTimeout = clusterManagerConfig.getContainerRequestTimeout(); + this.requestExpiryTimeout = clusterManagerConfig.getContainerRequestTimeout(); } /** @@ -160,20 +177,24 @@ public void run() { * During the run() method, the thread sleeps for allocatorSleepIntervalMs ms. It then invokes assignResourceRequests, * and tries to allocate any unsatisfied request that is still in the request queue {@link ResourceRequestState}) * with allocated resources. - * When {@code hostAffinityEnabled} is disabled, all allocated resources are buffered in the list keyed by "ANY_HOST". - * When {@code hostAffinityEnabled} is enabled, all allocated resources are buffered in the list keyed by "preferredHost + * When host-affinity is disabled, all allocated resources are buffered by the key "ANY_HOST". + * When host-affinity is disabled, all allocated resources are buffered by the hostName as key + * + * * * If the requested host is not available, the thread checks to see if the request has expired. If it has expired - * then two cases are handled seperately + * then two cases are handled separately * - * Case 1: host-affinity is disabled, cancels the current request and issues another ANY_HOST request - * Case 2: host-affinity is enabled, looks for allocated resouces on ANY_HOST and issues a container start if available, + * Case 1: host-affinity is enabled, looks for allocated resouces on ANY_HOST and issues a container start if available, * otherwise issues an ANY_HOST request + * Case 2: host-affinity is disabled, expired requests are not handled, allocator waits for cluster manager to issue + * resources + * TODO: SAMZA-2330 Hadle expired request for host affinity disabled case * - * In either of the scenarious if a {@code StandbyContainerManager} is present, the allocator transfers the request - * to it for checking StandByConstraints + * When host-affinity is enabled and a {@code StandbyContainerManager} is present, the allocator transfers the request + * to it for checking StandByConstraints before launcing a processor */ - protected void assignResourceRequests() { + void assignResourceRequests() { while (hasReadyPendingRequest()) { SamzaResourceRequest request = peekReadyPendingRequest().get(); String processorId = request.getProcessorId(); @@ -191,10 +212,9 @@ protected void assignResourceRequests() { state.matchedResourceRequests.incrementAndGet(); } - // Try to launch processor on this preferredHost if it all standby constraints are met + // If hot-standby is enabled, check standby constraints are met before launching a processor if (this.standbyContainerManager.isPresent()) { - standbyContainerManager.get().checkStandbyConstraintsAndRunStreamProcessor(request, preferredHost, - peekAllocatedResource(preferredHost), this, resourceRequestState); + checkStandByContrainsAndRunStreamProcessor(request, preferredHost); } else { runStreamProcessor(request, preferredHost); } @@ -213,7 +233,7 @@ protected void assignResourceRequests() { } else { LOG.info("Request for Processor ID: {} on preferred host {} has not expired yet." + "Request creation time: {}. Current Time: {}. Request timeout: {} ms", processorId, preferredHost, - requestCreationTime, System.currentTimeMillis(), requestTimeout); + requestCreationTime, System.currentTimeMillis(), requestExpiryTimeout); break; } } @@ -224,7 +244,8 @@ protected void assignResourceRequests() { * Handles an expired resource request when {@code hostAffinityEnabled} is true, in this case since the * preferred host, we try to see if a surplus ANY_HOST is available in the request queue. */ - private void handleExpiredRequestWithHostAffinityEnabled(String processorId, String preferredHost, + @VisibleForTesting + void handleExpiredRequestWithHostAffinityEnabled(String processorId, String preferredHost, SamzaResourceRequest request) { boolean resourceAvailableOnAnyHost = hasAllocatedResource(ResourceRequestState.ANY_HOST); if (standbyContainerManager.isPresent()) { @@ -276,6 +297,16 @@ protected void runStreamProcessor(SamzaResourceRequest request, String preferred clusterResourceManager.launchStreamProcessor(resource, builder); } + /** + * If {@code StandbyContainerManager} is present check standBy constraints are met before attempting to launch + * @param request outstanding request which has an allocated resource + * @param preferredHost to run the request + */ + private void checkStandByContrainsAndRunStreamProcessor(SamzaResourceRequest request, String preferredHost) { + standbyContainerManager.get().checkStandbyConstraintsAndRunStreamProcessor(request, preferredHost, + peekAllocatedResource(preferredHost), this, resourceRequestState); + } + /** * Called during initial request for resources * @@ -438,10 +469,10 @@ public void stop() { */ private boolean isRequestExpired(SamzaResourceRequest request) { long currTime = Instant.now().toEpochMilli(); - boolean requestExpired = currTime - request.getRequestTimestamp().toEpochMilli() > requestTimeout; + boolean requestExpired = currTime - request.getRequestTimestamp().toEpochMilli() > requestExpiryTimeout; if (requestExpired) { LOG.info("Request for Processor ID: {} on host: {} with creation time: {} has expired at current time: {} after timeout: {} ms.", - request.getProcessorId(), request.getPreferredHost(), request.getRequestTimestamp(), currTime, requestTimeout); + request.getProcessorId(), request.getPreferredHost(), request.getRequestTimestamp(), currTime, requestExpiryTimeout); } return requestExpired; } diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java index 7ccc27936f..ea8d6ef58e 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java @@ -22,10 +22,13 @@ import java.lang.reflect.Field; import java.time.Duration; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; import org.apache.samza.config.Config; import org.apache.samza.config.MapConfig; import org.apache.samza.container.LocalityManager; @@ -39,14 +42,17 @@ import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mockito; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.runners.MockitoJUnitRunner; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertNull; -import static org.junit.Assert.assertTrue; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; +import static org.junit.Assert.*; +import static org.mockito.Mockito.*; + +@RunWith(MockitoJUnitRunner.class) public class TestContainerAllocatorWithHostAffinity { private final Config config = getConfig(); @@ -69,9 +75,11 @@ private JobModelManager initializeJobModelManager(Config config, int containerCo } private AbstractContainerAllocator containerAllocator; + private AbstractContainerAllocator spyAllocator; private final int timeoutMillis = 1000; private MockContainerRequestState requestState; private Thread allocatorThread; + private Thread spyAllocatorThread; @Before public void setup() throws Exception { @@ -350,6 +358,147 @@ public void testExpiredRequestsAreCancelled() throws Exception { Assert.assertEquals(clusterResourceManager.cancelledRequests.size(), 3); } + @Test + public void testRequestAllocationOnPreferredHostWithRunStreamProcessor() throws Exception { + ClusterResourceManager.Callback mockCPM = mock(MockClusterResourceManagerCallback.class); + // Mock the callback from ClusterManager to add resources to the allocator + doAnswer((InvocationOnMock invocation) -> { + SamzaResource resource = (SamzaResource) invocation.getArgumentAt(0, List.class).get(0); + spyAllocator.addResource(resource); + return null; + }).when(mockCPM).onResourcesAvailable(anyList()); + + spyAllocator = Mockito.spy( + new AbstractContainerAllocator(new MockClusterResourceManager(mockCPM, state), config, state, + getClass().getClassLoader(), true, Optional.empty())); + + // Request Resources + spyAllocator.requestResources(new HashMap() { + { + put("0", "abc"); + put("1", "xyz"); + } + }); + + spyAllocatorThread = new Thread(spyAllocator); + + // Start the container allocator thread periodic assignment + spyAllocatorThread.start(); + // Let Allocator thread periodically fulfill requests + Thread.sleep(100); + + // Verify that all the request that were created were preferred host requests + ArgumentCaptor resourceRequestCaptor = ArgumentCaptor.forClass(SamzaResourceRequest.class); + verify(spyAllocator, times(2)).runStreamProcessor(resourceRequestCaptor.capture(), anyString()); + resourceRequestCaptor.getAllValues() + .forEach(resourceRequest -> assertNotEquals(resourceRequest.getPreferredHost(), ResourceRequestState.ANY_HOST)); + Set hostNames = resourceRequestCaptor.getAllValues().stream().map(request -> request.getPreferredHost()).collect( + Collectors.toSet()); + assertTrue(hostNames.contains("abc")); + assertTrue(hostNames.contains("xyz")); + // No any host requests should be made if preferred host is satisfied + assertTrue(state.anyHostRequests.get() == 0); + // State check when host affinity is enabled + assertTrue(state.matchedResourceRequests.get() == 2); + assertTrue(state.preferredHostRequests.get() == 2); + containerAllocator.stop(); + } + + @Test + public void testExpiredRequestAllocationOnAnyHost() throws Exception { + MockClusterResourceManager spyManager = spy(new MockClusterResourceManager(callback, state)); + spyAllocator = Mockito.spy( + new AbstractContainerAllocator(spyManager, config, state, + getClass().getClassLoader(), true, Optional.empty())); + + // Request Preferred Resources + spyAllocator.requestResources(new HashMap() { + { + put("0", "abc"); + put("1", "def"); + } + }); + + spyAllocatorThread = new Thread(spyAllocator); + // Start the container allocator thread periodic assignment + spyAllocatorThread.start(); + + // Let the request expire, expiration timeout is 3 ms + Thread.sleep(10); + + // Verify that all the request that were created as preferred host requests expired + assertTrue(state.preferredHostRequests.get() == 2); + assertTrue(state.expiredPreferredHostRequests.get() == 2); + verify(spyAllocator, times(1)).handleExpiredRequestWithHostAffinityEnabled(eq("0"), eq("abc"), + any(SamzaResourceRequest.class)); + verify(spyAllocator, times(1)).handleExpiredRequestWithHostAffinityEnabled(eq("1"), eq("def"), + any(SamzaResourceRequest.class)); + + // Verify that preferred host request were cancelled and since no surplus resources were available + // requestResource was invoked with ANY_HOST requests + ArgumentCaptor cancelledRequestCaptor = ArgumentCaptor.forClass(SamzaResourceRequest.class); + // At least 2 preferred host requests were cancelled + verify(spyManager, atLeast(2)).cancelResourceRequest(cancelledRequestCaptor.capture()); + assertTrue(cancelledRequestCaptor.getAllValues().stream().map(resourceRequest -> resourceRequest.getPreferredHost()).collect( + Collectors.toSet()).size() > 2); + // Check that atleast 2 ANY_HOST requests were made + assertTrue(state.matchedResourceRequests.get() == 0); + assertTrue(state.anyHostRequests.get() > 2); + containerAllocator.stop(); + } + + @Test + public void testExpiredRequestAllocationOnSurplusAnyHostWithRunStreamProcessor() throws Exception { + // Add Extra Resources + spyAllocator = Mockito.spy( + new AbstractContainerAllocator(new MockClusterResourceManager(callback, state), config, state, + getClass().getClassLoader(), true, Optional.empty())); + + spyAllocator.addResource(new SamzaResource(1, 1000, "xyz", "id1")); + spyAllocator.addResource(new SamzaResource(1, 1000, "zzz", "id2")); + + // Request Preferred Resources + spyAllocator.requestResources(new HashMap() { + { + put("0", "abc"); + put("1", "def"); + } + }); + + spyAllocatorThread = new Thread(spyAllocator); + // Start the container allocator thread periodic assignment + spyAllocatorThread.start(); + + // Let the request expire, expiration timeout is 3 ms + Thread.sleep(10); + + // Verify that all the request that were created as preferred host requests expired + assertTrue(state.expiredPreferredHostRequests.get() == 2); + verify(spyAllocator, times(1)).handleExpiredRequestWithHostAffinityEnabled(eq("0"), eq("abc"), + any(SamzaResourceRequest.class)); + verify(spyAllocator, times(1)).handleExpiredRequestWithHostAffinityEnabled(eq("1"), eq("def"), + any(SamzaResourceRequest.class)); + + // Verify that runStreamProcessor was invoked with already available ANY_HOST requests + ArgumentCaptor resourceRequestCaptor = ArgumentCaptor.forClass(SamzaResourceRequest.class); + ArgumentCaptor hostCaptor = ArgumentCaptor.forClass(String.class); + verify(spyAllocator, times(2)).runStreamProcessor(resourceRequestCaptor.capture(), hostCaptor.capture()); + // Resource request were preferred host requests + resourceRequestCaptor.getAllValues() + .forEach(resourceRequest -> assertNotEquals(resourceRequest.getPreferredHost(), ResourceRequestState.ANY_HOST)); + + // Since requests expired, allocator ran the requests on surplus available ANY_HOST + hostCaptor.getAllValues() + .forEach(host -> assertEquals(host, ResourceRequestState.ANY_HOST)); + + // State Update check + assertTrue(state.matchedResourceRequests.get() == 0); + assertTrue(state.preferredHostRequests.get() == 2); + assertTrue(state.anyHostRequests.get() == 0); + containerAllocator.stop(); + } + + //@Test public void testExpiryWithNonResponsiveClusterManager() throws Exception { diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java index 978055a657..52f95d35fc 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java @@ -21,6 +21,7 @@ import java.lang.reflect.Field; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Optional; import org.apache.samza.config.Config; @@ -33,12 +34,21 @@ import org.junit.After; import org.junit.Before; import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mockito; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.runners.MockitoJUnitRunner; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; +import static org.mockito.Matchers.*; +import static org.mockito.Mockito.*; + +@RunWith(MockitoJUnitRunner.class) public class TestContainerAllocatorWithoutHostAffinity { private final MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); private final Config config = getConfig(); @@ -52,6 +62,9 @@ public class TestContainerAllocatorWithoutHostAffinity { private MockContainerRequestState requestState; private Thread allocatorThread; + private Thread spyThread; + private AbstractContainerAllocator spyAllocator; + @Before public void setup() throws Exception { containerAllocator = new AbstractContainerAllocator(manager, config, state, getClass().getClassLoader(), false, Optional @@ -66,7 +79,7 @@ public void setup() throws Exception { @After public void teardown() throws Exception { jobModelManager.stop(); - containerAllocator.stop(); + validateMockitoUsage(); } private static Config getConfig() { @@ -101,7 +114,6 @@ public void testAddContainer() throws Exception { containerAllocator.addResource(new SamzaResource(1, 1000, "abc", "id1")); containerAllocator.addResource(new SamzaResource(1, 1000, "xyz", "id1")); - assertNull(requestState.getResourcesOnAHost("abc")); assertNotNull(requestState.getResourcesOnAHost(ResourceRequestState.ANY_HOST)); assertTrue(requestState.getResourcesOnAHost(ResourceRequestState.ANY_HOST).size() == 2); @@ -235,4 +247,49 @@ public void run() { listener.verify(); } + + /** + * Test the complete flow from container request creation to allocation when host affinity is disabled + */ + @Test + public void testRequestAllocationWithRunStreamProcessor() throws Exception { + Map containersToHostMapping = new HashMap() { + { + put("0", "prev_host"); + put("1", "prev_host"); + put("2", "prev_host"); + put("3", "prev_host"); + } + }; + + ClusterResourceManager.Callback mockCPM = mock(ClusterResourceManager.Callback.class); + spyAllocator = Mockito.spy( + new AbstractContainerAllocator(new MockClusterResourceManager(mockCPM, state), config, state, + getClass().getClassLoader(), false, Optional.empty())); + // Mock the callback from ClusterManager to add resources to the allocator + doAnswer((InvocationOnMock invocation) -> { + SamzaResource resource = (SamzaResource) invocation.getArgumentAt(0, List.class).get(0); + spyAllocator.addResource(resource); + return null; + }).when(mockCPM).onResourcesAvailable(anyList()); + // Request Resources + spyAllocator.requestResources(containersToHostMapping); + spyThread = new Thread(spyAllocator, "Container Allocator Thread"); + // Start the container allocator thread periodic assignment + spyThread.start(); + Thread.sleep(1000); + // Verify that all the request that were created were "ANY_HOST" requests + ArgumentCaptor resourceRequestCaptor = ArgumentCaptor.forClass(SamzaResourceRequest.class); + verify(spyAllocator, times(4)).runStreamProcessor(resourceRequestCaptor.capture(), anyString()); + resourceRequestCaptor.getAllValues() + .forEach(resourceRequest -> assertEquals(resourceRequest.getPreferredHost(), ResourceRequestState.ANY_HOST)); + assertTrue(state.anyHostRequests.get() == containersToHostMapping.size()); + // Expiry currently is only handled for host affinity enabled cases + verify(spyAllocator, never()).handleExpiredRequestWithHostAffinityEnabled(anyString(), anyString(), + any(SamzaResourceRequest.class)); + // Only updated when host affinity is enabled + assertTrue(state.matchedResourceRequests.get() == 0); + assertTrue(state.preferredHostRequests.get() == 0); + spyAllocator.stop(); + } } diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java index 7996e59944..8929106bd7 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerProcessManager.java @@ -71,7 +71,7 @@ public class TestContainerProcessManager { put("cluster-manager.container.retry.count", "1"); put("cluster-manager.container.retry.window.ms", "1999999999"); put("cluster-manager.allocator.sleep.ms", "1"); - put("cluster-manager.container.request.timeout.ms", "5000"); + put("cluster-manager.container.request.timeout.ms", "2"); put("cluster-manager.container.memory.mb", "512"); put("yarn.package.path", "/foo"); put("task.inputs", "test-system.test-stream"); @@ -762,7 +762,9 @@ public void testDuplicateNotificationsDoNotAffectJobHealth() throws Exception { @Test public void testNewContainerRequestedOnFailureWithKnownCode() throws Exception { Config conf = getConfig(); - Map config = new HashMap<>(getConfig()); + + Map config = new HashMap<>(); + config.putAll(getConfig()); SamzaApplicationState state = new SamzaApplicationState(getJobModelManagerWithoutHostAffinity(1)); MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback(); MockClusterResourceManager clusterResourceManager = new MockClusterResourceManager(callback, state); diff --git a/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java b/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java index d1acdd2c07..57c2c34cbd 100644 --- a/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java +++ b/samza-core/src/test/java/org/apache/samza/diagnostics/TestDiagnosticsManager.java @@ -71,7 +71,8 @@ public void setup() { Mockito.when(mockExecutorService.scheduleWithFixedDelay(Mockito.any(), Mockito.anyLong(), Mockito.anyLong(), Mockito.eq(TimeUnit.SECONDS))).thenAnswer(invocation -> { ((Runnable) invocation.getArguments()[0]).run(); - return Mockito.mock(ScheduledFuture.class); + return Mockito. + mock(ScheduledFuture.class); }); this.diagnosticsManager = From 10e3b683ad1a737f5994c70e9b1b53f67ef4057b Mon Sep 17 00:00:00 2001 From: Sanil15 Date: Tue, 24 Sep 2019 10:51:25 -0700 Subject: [PATCH 4/6] Removing whitespace --- .../samza/clustermanager/AbstractContainerAllocator.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java index c530d05af2..4ce87de43c 100644 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java +++ b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java @@ -177,11 +177,10 @@ public void run() { * During the run() method, the thread sleeps for allocatorSleepIntervalMs ms. It then invokes assignResourceRequests, * and tries to allocate any unsatisfied request that is still in the request queue {@link ResourceRequestState}) * with allocated resources. + * * When host-affinity is disabled, all allocated resources are buffered by the key "ANY_HOST". * When host-affinity is disabled, all allocated resources are buffered by the hostName as key * - * - * * If the requested host is not available, the thread checks to see if the request has expired. If it has expired * then two cases are handled separately * From e4b565f25433e05895a9b240d2ca9308e835108b Mon Sep 17 00:00:00 2001 From: Sanil15 Date: Tue, 24 Sep 2019 10:54:45 -0700 Subject: [PATCH 5/6] Removing whitespace --- .../clustermanager/TestContainerAllocatorWithHostAffinity.java | 1 - .../TestContainerAllocatorWithoutHostAffinity.java | 1 - 2 files changed, 2 deletions(-) diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java index ea8d6ef58e..4ae63b214b 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithHostAffinity.java @@ -498,7 +498,6 @@ public void testExpiredRequestAllocationOnSurplusAnyHostWithRunStreamProcessor() containerAllocator.stop(); } - //@Test public void testExpiryWithNonResponsiveClusterManager() throws Exception { diff --git a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java index 52f95d35fc..e75eff35b7 100644 --- a/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java +++ b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java @@ -247,7 +247,6 @@ public void run() { listener.verify(); } - /** * Test the complete flow from container request creation to allocation when host affinity is disabled */ From ceb0fda5c20f2c43e538e61c93496baa7fed27db Mon Sep 17 00:00:00 2001 From: Sanil15 Date: Fri, 27 Sep 2019 12:13:01 -0700 Subject: [PATCH 6/6] Addressing Review, updating docs --- .../samza/clustermanager/AbstractContainerAllocator.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java index 4ce87de43c..6a5eb7a628 100644 --- a/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java +++ b/samza-core/src/main/java/org/apache/samza/clustermanager/AbstractContainerAllocator.java @@ -45,15 +45,15 @@ * in the coordinator stream * *
  • - * This thread periodically periodically matches outstanding resource requests with allocated resources. + * This thread periodically matches outstanding resource requests with allocated resources. * Its period is controlled using the {@code allocatorSleepIntervalMs} parameter *
  • *
  • - * In case of host-affinity is enabled, the resource-request's preferredHost param is set to the host the processor + * When host-affinity is enabled, the resource-request's preferredHost param is set to the host the processor * was last seen on *
  • *
  • - * In case of host-affinity us disabled, the resource-request's preferredHost param is set to {@link ResourceRequestState#ANY_HOST} + * When host-affinity is disabled, the resource-request's preferredHost param is set to {@link ResourceRequestState#ANY_HOST} *
  • *
  • * When host-affinity is enabled and a preferred resource has not been obtained after {@code requestExpiryTimeout} @@ -179,7 +179,7 @@ public void run() { * with allocated resources. * * When host-affinity is disabled, all allocated resources are buffered by the key "ANY_HOST". - * When host-affinity is disabled, all allocated resources are buffered by the hostName as key + * When host-affinity is enabled, all allocated resources are buffered by the hostName as key * * If the requested host is not available, the thread checks to see if the request has expired. If it has expired * then two cases are handled separately