Github user tillrohrmann commented on a diff in the pull request:
https://github.com/apache/flink/pull/4628#discussion_r140032946
--- Diff:
flink-mesos/src/main/java/org/apache/flink/mesos/runtime/clusterframework/MesosFlinkResourceManager.java
---
@@ -675,6 +679,42 @@ private LaunchableMesosWorker
createLaunchableMesosWorker(Protos.TaskID taskID)
}
/**
+ * Sets a coTaskGetter callback for evaluating balancing constraint.
+ */
+ private void setCoTaskGetter() {
+ for
(MesosTaskManagerParameters.BalancedHostAttrConstraintParams param :
taskManagerParameters.balancedConstraintParams()) {
+ param.setCoTasksGetter(new Func1<String, Set<String>>()
{
+ @Override
+ public Set<String> call(String s) {
+ Map<String, Set<String>>
taskToCoTasksMap = new HashMap<>();
+ Set <String> taskIds = getTaskIdsSet();
+ for (String taskId : taskIds) {
+ Set <String> 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<String> getTaskIdsSet() {
+ Set<String> taskIds = new HashSet<String>();
+ List <MesosWorkerStore.Worker> workers = new
ArrayList<MesosWorkerStore.Worker>();
+ workers.addAll(this.workersInNew.values());
+ workers.addAll(this.workersInLaunch.values());
--- End diff --
I think we don't have to add the `workersInNew.values()` and
`workersInLaunch.values()` first to `workers` and then only to `taskIds`. We
can directly add them to `taskIds`. Saves us one copy operation.
---