From c333344ba0b57d94bb3e1ae7c88e50d2a8549a8b Mon Sep 17 00:00:00 2001 From: bbayani Date: Wed, 30 Aug 2017 22:06:45 +0530 Subject: [PATCH] [FLINK-7486]:[flink-mesos]:Support for adding unique attribute / group_by attribute constraints --- docs/ops/config.md | 2 + docs/ops/deployment/mesos.md | 4 + .../LaunchableMesosWorker.java | 2 +- .../MesosFlinkResourceManager.java | 40 +++++++++ .../MesosResourceManager.java | 44 ++++++++++ .../MesosTaskManagerParameters.java | 87 +++++++++++++++++++ .../MesosFlinkResourceManagerTest.java | 1 + .../MesosResourceManagerTest.java | 5 +- 8 files changed, 182 insertions(+), 3 deletions(-) diff --git a/docs/ops/config.md b/docs/ops/config.md index e0b9d4db714f50..2e1f51b057de64 100644 --- a/docs/ops/config.md +++ b/docs/ops/config.md @@ -477,6 +477,8 @@ use the `env.java.opts` setting, which is the `%jvmopts%` variable in the String - `mesos.constraints.hard.hostattribute`: Constraints for task placement on mesos (**DEFAULT**: None). +- `mesos.constraints.soft.balanced`: Soft Constraints for balancing the tasks across mesos based on agent attributes (**DEFAULT**: None). + - `mesos.maximum-failed-tasks`: The maximum number of failed workers before the cluster fails (**DEFAULT**: Number of initial workers). May be set to -1 to disable this feature. diff --git a/docs/ops/deployment/mesos.md b/docs/ops/deployment/mesos.md index 2fa340d6d6d30b..356e5a820c340d 100644 --- a/docs/ops/deployment/mesos.md +++ b/docs/ops/deployment/mesos.md @@ -226,6 +226,10 @@ When running Flink with Marathon, the whole Flink cluster including the job mana Takes a comma-separated list of key:value pairs corresponding to the attributes exposed by the target mesos agents. Example: `az:eu-west-1a,series:t2` +`mesos.constraints.soft.balanced`: Soft Constraints for balancing the tasks across mesos based on agent attributes (**DEFAULT**: None). +Takes a comma-separated list of key=value pairs. Key corresponds to host attribute and value is number of expected unique values for given host attribute. +Example: `az=3,rack_id=4` + `mesos.maximum-failed-tasks`: The maximum number of failed workers before the cluster fails (**DEFAULT**: Number of initial workers). May be set to -1 to disable this feature. diff --git a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/LaunchableMesosWorker.java b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/LaunchableMesosWorker.java index 2c3250738027c7..da0dadb91c3506 100644 --- a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/LaunchableMesosWorker.java +++ b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/LaunchableMesosWorker.java @@ -155,7 +155,7 @@ public List getHardConstraints() { @Override public List getSoftConstraints() { - return null; + return params.softConstraints(); } @Override diff --git a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManager.java b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManager.java index 6335745004a18b..39e751aabdd08a 100644 --- a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManager.java +++ b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManager.java @@ -55,6 +55,7 @@ import com.netflix.fenzo.TaskScheduler; import com.netflix.fenzo.VirtualMachineLease; import com.netflix.fenzo.functions.Action1; +import com.netflix.fenzo.functions.Func1; import org.apache.mesos.Protos; import org.apache.mesos.Protos.FrameworkInfo; import org.apache.mesos.SchedulerDriver; @@ -63,8 +64,10 @@ import java.util.ArrayList; import java.util.Collection; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; import scala.Option; @@ -663,6 +666,7 @@ private void taskTerminated(Protos.TaskID taskID, Protos.TaskStatus status) { // ------------------------------------------------------------------------ private LaunchableMesosWorker createLaunchableMesosWorker(Protos.TaskID taskID) { + setCoTaskGetter(); LaunchableMesosWorker launchable = new LaunchableMesosWorker( artifactResolver, @@ -674,6 +678,42 @@ private LaunchableMesosWorker createLaunchableMesosWorker(Protos.TaskID taskID) return launchable; } + /** + * Sets a coTaskGetter callback for evaluating balancing constraint. + */ + private void setCoTaskGetter() { + for (MesosTaskManagerParameters.BalancedHostAttrConstraintParams param : taskManagerParameters.balancedConstraintParams()) { + param.setCoTasksGetter(new Func1>() { + @Override + public Set call(String s) { + Map> taskToCoTasksMap = new HashMap<>(); + Set taskIds = getTaskIdsSet(); + for (String taskId : taskIds) { + Set coTaskIds = new HashSet<>(taskIds); + coTaskIds.remove(taskId); + taskToCoTasksMap.put(taskId, coTaskIds); + } + return taskToCoTasksMap.get(s); + } + }); + } + } + + /** + * Compiles the set of task IDs in new/launch state. + * @return The unique TaskIDs + */ + private Set getTaskIdsSet() { + Set taskIds = new HashSet(); + List workers = new ArrayList(); + workers.addAll(this.workersInNew.values()); + workers.addAll(this.workersInLaunch.values()); + for (MesosWorkerStore.Worker worker : workers) { + taskIds.add(worker.taskID().getValue()); + } + return taskIds; + } + /** * Extracts a unique ResourceID from the Mesos task. * diff --git a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManager.java b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManager.java index 9a2ad42a48cb55..882d98a3be9052 100644 --- a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManager.java +++ b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManager.java @@ -69,6 +69,7 @@ import com.netflix.fenzo.TaskScheduler; import com.netflix.fenzo.VirtualMachineLease; import com.netflix.fenzo.functions.Action1; +import com.netflix.fenzo.functions.Func1; import org.apache.mesos.Protos; import org.apache.mesos.Scheduler; import org.apache.mesos.SchedulerDriver; @@ -79,8 +80,10 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; @@ -669,6 +672,7 @@ private LaunchableMesosWorker createLaunchableMesosWorker(Protos.TaskID taskID, new HashMap<>(taskManagerParameters.containeredParameters().taskManagerEnv())), taskManagerParameters.containerVolumes(), taskManagerParameters.constraints(), + taskManagerParameters.balancedConstraintParams(), taskManagerParameters.command(), taskManagerParameters.bootstrapCommand(), taskManagerParameters.getTaskManagerHostname() @@ -676,6 +680,8 @@ private LaunchableMesosWorker createLaunchableMesosWorker(Protos.TaskID taskID, LOG.debug("LaunchableMesosWorker parameters: {}", params); + setCoTaskGetter(); + LaunchableMesosWorker launchable = new LaunchableMesosWorker( artifactServer, @@ -687,6 +693,44 @@ private LaunchableMesosWorker createLaunchableMesosWorker(Protos.TaskID taskID, return launchable; } + /** + * Sets a coTaskGetter callback for evaluating balancing constraint. + */ + private void setCoTaskGetter() { + for (MesosTaskManagerParameters.BalancedHostAttrConstraintParams param : taskManagerParameters.balancedConstraintParams()) { + param.setCoTasksGetter(new Func1>() { + @Override + public Set call(String s) { + Map> taskToCoTasksMap = new HashMap<>(); + Set taskIds = getTaskIdsSet(); + for (String taskId : taskIds) { + Set coTaskIds = new HashSet<>(taskIds); + coTaskIds.remove(taskId); + taskToCoTasksMap.put(taskId, coTaskIds); + } + return taskToCoTasksMap.get(s); + } + }); + } + } + + + /** + * Compiles the set of task IDs in new/launch state. + * @return The unique TaskIDs + */ + + private Set getTaskIdsSet() { + Set taskIds = new HashSet(); + List workers = new ArrayList(); + workers.addAll(this.workersInNew.values()); + workers.addAll(this.workersInLaunch.values()); + for (MesosWorkerStore.Worker worker : workers) { + taskIds.add(worker.taskID().getValue()); + } + return taskIds; + } + /** * Extracts a unique ResourceID from the Mesos task. * diff --git a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosTaskManagerParameters.java b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosTaskManagerParameters.java index 3859913ecda3d3..e37115b0ec0137 100644 --- a/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosTaskManagerParameters.java +++ b/flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosTaskManagerParameters.java @@ -26,13 +26,17 @@ import org.apache.flink.util.Preconditions; import com.netflix.fenzo.ConstraintEvaluator; +import com.netflix.fenzo.VMTaskFitnessCalculator; import com.netflix.fenzo.functions.Func1; +import com.netflix.fenzo.plugins.BalancedHostAttrConstraint; import com.netflix.fenzo.plugins.HostAttrValueConstraint; + import org.apache.mesos.Protos; import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Set; import java.util.regex.Pattern; import scala.Option; @@ -90,6 +94,10 @@ public class MesosTaskManagerParameters { key("mesos.constraints.hard.hostattribute") .noDefaultValue(); + public static final ConfigOption MESOS_CONSTRAINTS_SOFT_BALANCED = + key("mesos.constraints.soft.balanced") + .noDefaultValue(); + /** * Value for {@code MESOS_RESOURCEMANAGER_TASKS_CONTAINER_TYPE} setting. Tells to use the Mesos containerizer. */ @@ -111,6 +119,10 @@ public class MesosTaskManagerParameters { private final List constraints; + private List softConstraints; + + private final List balancedConstraintParams; + private final String command; private final Option bootstrapCommand; @@ -124,6 +136,7 @@ public MesosTaskManagerParameters( ContaineredTaskManagerParameters containeredParameters, List containerVolumes, List constraints, + List balancedConstraintParams, String command, Option bootstrapCommand, Option taskManagerHostname) { @@ -134,6 +147,7 @@ public MesosTaskManagerParameters( this.containeredParameters = Preconditions.checkNotNull(containeredParameters); this.containerVolumes = Preconditions.checkNotNull(containerVolumes); this.constraints = Preconditions.checkNotNull(constraints); + this.balancedConstraintParams = Preconditions.checkNotNull(balancedConstraintParams); this.command = Preconditions.checkNotNull(command); this.bootstrapCommand = Preconditions.checkNotNull(bootstrapCommand); this.taskManagerHostname = Preconditions.checkNotNull(taskManagerHostname); @@ -183,6 +197,25 @@ public List constraints() { return constraints; } + /** + * Get the balanced constraints parameters. + */ + public List balancedConstraintParams() { + return balancedConstraintParams; + } + + /** + * Get the placement soft constraints. + */ + public List softConstraints() { + List softConstraints = new ArrayList<>(); + for (MesosTaskManagerParameters.BalancedHostAttrConstraintParams param : balancedConstraintParams) { + BalancedHostAttrConstraint aConstraint = new BalancedHostAttrConstraint(param.coTasksGetter, param.hostAttr, Integer.parseInt(param.numOfExpectedUniqueValues)); + softConstraints.add(aConstraint.asSoftConstraint()); + } + return softConstraints; + } + /** * Get the taskManager hostname. */ @@ -213,6 +246,7 @@ public String toString() { ", containeredParameters=" + containeredParameters + ", containerVolumes=" + containerVolumes + ", constraints=" + constraints + + ", softConstraints=" + softConstraints + ", taskManagerHostName=" + taskManagerHostname + ", command=" + command + ", bootstrapCommand=" + bootstrapCommand + @@ -227,6 +261,8 @@ public String toString() { public static MesosTaskManagerParameters create(Configuration flinkConfig) { List constraints = parseConstraints(flinkConfig.getString(MESOS_CONSTRAINTS_HARD_HOSTATTR)); + List balancedConstraintParams = + parseSoftConstraints(flinkConfig.getString(MESOS_CONSTRAINTS_SOFT_BALANCED)); // parse the common parameters ContaineredTaskManagerParameters containeredParameters = ContaineredTaskManagerParameters.create( flinkConfig, @@ -276,6 +312,7 @@ public static MesosTaskManagerParameters create(Configuration flinkConfig) { containeredParameters, containerVolumes, constraints, + balancedConstraintParams, tmCommand, tmBootstrapCommand, taskManagerHostname); @@ -312,6 +349,28 @@ public String call(String s) { })); } + private static List parseSoftConstraints(String mesosConstraints) { + + if (mesosConstraints == null || mesosConstraints.isEmpty()) { + return Collections.emptyList(); + } else { + List constraints = new ArrayList<>(); + + for (String constraint : mesosConstraints.split(",")) { + if (constraint.isEmpty()) { + continue; + } + final String[] constraintList = constraint.split("="); + if (constraintList.length != 2) { + continue; + } + constraints.add(new MesosTaskManagerParameters.BalancedHostAttrConstraintParams(constraintList[0], constraintList[1])); + } + + return constraints; + } + } + /** * Used to build volume specs for mesos. This allows for mounting additional volumes into a container * @@ -365,6 +424,34 @@ public static List buildVolumes(Option containerVolumes) } } + /** + * Internal class that stores the parsed information about soft constraint + * It encapsulates fields: + * 1. Host attribute name + * 2. Expected number of unique values for given host attribute + * 3. A callback coTaskGetter used while evaluating balancing constraint + */ + static class BalancedHostAttrConstraintParams { + String hostAttr; + String numOfExpectedUniqueValues; + + Func1> coTasksGetter; + + public BalancedHostAttrConstraintParams(String hostAttr, String numOfExpectedUniqueValues) { + this.hostAttr = hostAttr; + this.numOfExpectedUniqueValues = numOfExpectedUniqueValues; + } + + public void setCoTasksGetter(Func1> coTasksGetter) { + this.coTasksGetter = coTasksGetter; + } + + @Override + public String toString() { + return "{ Host Attribute: " + this.hostAttr + ", Number of expected unique values: " + this.numOfExpectedUniqueValues + "}"; + } + } + /** * The supported containerizers. */ diff --git a/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManagerTest.java b/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManagerTest.java index ff324865274e7c..859ba02a9a5331 100644 --- a/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManagerTest.java +++ b/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManagerTest.java @@ -251,6 +251,7 @@ public void initialize() { containeredParams, Collections.emptyList(), Collections.emptyList(), + Collections.emptyList(), "", Option.empty(), Option.empty()); diff --git a/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManagerTest.java b/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManagerTest.java index cf0c9135364751..bfb772523e1cef 100644 --- a/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManagerTest.java +++ b/flink-mesos/src/test/java/org/apache/flink/mesos/runtime/clusterframework/MesosResourceManagerTest.java @@ -251,8 +251,9 @@ static class Context implements AutoCloseable { new ContaineredTaskManagerParameters(1024, 768, 256, 4, new HashMap()); MesosTaskManagerParameters tmParams = new MesosTaskManagerParameters( 1.0, MesosTaskManagerParameters.ContainerType.MESOS, Option.empty(), containeredParams, - Collections.emptyList(), Collections.emptyList(), "", Option.empty(), - Option.empty()); + Collections.emptyList(), Collections.emptyList(), + Collections.emptyList(), + "", Option.empty(), Option.empty()); // resource manager rmConfiguration = new ResourceManagerConfiguration(