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..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
@@ -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;
@@ -35,20 +36,51 @@
/**
* {@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} 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 matches outstanding resource requests with allocated resources.
+ * Its period is controlled using the {@code allocatorSleepIntervalMs} parameter
+ *
+ * -
+ * When host-affinity is enabled, the resource-request's preferredHost param is set to the host the processor
+ * was last seen on
+ *
+ * -
+ * 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}
+ * 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
*/
-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 +115,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 requestExpiryTimeout;
- 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.requestExpiryTimeout = clusterManagerConfig.getContainerRequestTimeout();
}
/**
@@ -121,18 +163,104 @@ 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 host-affinity is disabled, all allocated resources are buffered by the key "ANY_HOST".
+ * 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
+ *
+ * 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
+ *
+ * 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 abstract void assignResourceRequests();
+ 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();
+ }
+
+ // If hot-standby is enabled, check standby constraints are met before launching a processor
+ if (this.standbyContainerManager.isPresent()) {
+ checkStandByContrainsAndRunStreamProcessor(request, preferredHost);
+ } 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 {
+ 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(), requestExpiryTimeout);
+ 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.
+ */
+ @VisibleForTesting
+ 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);
+ }
+ }
/**
* Updates the request state and runs a processor on the specified host. Assumes a resource
@@ -156,7 +284,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
@@ -168,6 +296,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
*
@@ -178,7 +316,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 +459,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() > 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, requestExpiryTimeout);
+ }
+ 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/MockHostAwareContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/MockContainerAllocatorWithHostAffinity.java
similarity index 87%
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 0c44edf2ee..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,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 MockContainerAllocatorWithHostAffinity extends AbstractContainerAllocator {
private Semaphore semaphore = new Semaphore(0);
- public MockHostAwareContainerAllocator(ClusterResourceManager manager,
+ public MockContainerAllocatorWithHostAffinity(ClusterResourceManager manager,
Config config, SamzaApplicationState state) {
- super(manager, ALLOCATOR_TIMEOUT_MS, config, Optional.empty(), state,
- MockHostAwareContainerAllocator.class.getClassLoader());
+ super(manager, config, state,
+ 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 88%
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 e294cfc5ad..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
@@ -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 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());
+ 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 72%
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 40d292c017..4ae63b214b 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
@@ -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,15 +42,18 @@
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.*;
-public class TestHostAwareContainerAllocator {
+
+@RunWith(MockitoJUnitRunner.class)
+public class TestContainerAllocatorWithHostAffinity {
private final Config config = getConfig();
private final JobModelManager jobModelManager = initializeJobModelManager(config, 1);
@@ -68,18 +74,20 @@ private JobModelManager initializeJobModelManager(Config config, int containerCo
new MockHttpServer("/", 7777, null, new ServletHolder(DefaultServlet.class)));
}
- private HostAwareContainerAllocator containerAllocator;
+ 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 {
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);
@@ -350,6 +358,146 @@ 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/TestContainerAllocator.java b/samza-core/src/test/java/org/apache/samza/clustermanager/TestContainerAllocatorWithoutHostAffinity.java
similarity index 72%
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 0f90f92b80..e75eff35b7 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
@@ -21,7 +21,9 @@
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;
import org.apache.samza.config.MapConfig;
import org.apache.samza.coordinator.JobModelManager;
@@ -32,13 +34,22 @@
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.*;
-public class TestContainerAllocator {
+
+@RunWith(MockitoJUnitRunner.class)
+public class TestContainerAllocatorWithoutHostAffinity {
private final MockClusterResourceManagerCallback callback = new MockClusterResourceManagerCallback();
private final Config config = getConfig();
private final JobModelManager jobModelManager = JobModelManagerTestUtil.getJobModelManager(config, 1,
@@ -47,15 +58,19 @@ 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;
+ private Thread spyThread;
+ private AbstractContainerAllocator spyAllocator;
+
@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);
@@ -64,7 +79,7 @@ public void setup() throws Exception {
@After
public void teardown() throws Exception {
jobModelManager.stop();
- containerAllocator.stop();
+ validateMockitoUsage();
}
private static Config getConfig() {
@@ -99,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);
@@ -233,4 +247,48 @@ 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 997adb1c5d..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
@@ -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);
@@ -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);
@@ -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",
@@ -637,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);
@@ -698,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);
@@ -769,7 +770,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);
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 =