Github user tillrohrmann commented on a diff in the pull request:
https://github.com/apache/flink/pull/4628#discussion_r140034034
--- 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);
--- End diff --
I think you're doing a lot of redundant work here. Wouldn't
`taskIds.remove(s)` simply do the same?
---