-
Notifications
You must be signed in to change notification settings - Fork 3.7k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
11 changed files
with
400 additions
and
493 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
42 changes: 0 additions & 42 deletions
42
...es-overlord-extensions/src/main/java/org/apache/druid/k8s/overlord/execution/Matcher.java
This file was deleted.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
157 changes: 157 additions & 0 deletions
157
...s-overlord-extensions/src/main/java/org/apache/druid/k8s/overlord/execution/Selector.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,157 @@ | ||
/* | ||
* 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.druid.k8s.overlord.execution; | ||
|
||
import com.fasterxml.jackson.annotation.JsonCreator; | ||
import com.fasterxml.jackson.annotation.JsonProperty; | ||
import org.apache.druid.indexing.common.task.Task; | ||
import org.apache.druid.query.DruidMetrics; | ||
|
||
import java.util.Map; | ||
import java.util.Objects; | ||
import java.util.Set; | ||
|
||
/** | ||
* Represents a condition-based selector that evaluates whether a given task meets specified criteria. | ||
* The selector uses conditions defined on context tags and task fields to determine if a task matches. | ||
*/ | ||
public class Selector | ||
{ | ||
private final String selectionKey; | ||
private final Map<String, Set<String>> cxtTagsConditions; | ||
private final Set<String> taskTypeCondition; | ||
private final Set<String> dataSourceCondition; | ||
|
||
/** | ||
* Creates a selector with specified conditions for context tags and task fields. | ||
* | ||
* @param selectionKey the identifier representing the outcome when a task matches the conditions | ||
* @param cxtTagsConditions conditions on context tags | ||
* @param taskTypeCondition conditions on task type | ||
* @param dataSourceCondition conditions on task dataSource | ||
*/ | ||
@JsonCreator | ||
public Selector( | ||
@JsonProperty("selectionKey") String selectionKey, | ||
@JsonProperty("context.tags") Map<String, Set<String>> cxtTagsConditions, | ||
@JsonProperty("type") Set<String> taskTypeCondition, | ||
@JsonProperty("dataSource") Set<String> dataSourceCondition | ||
) | ||
{ | ||
this.selectionKey = selectionKey; | ||
this.cxtTagsConditions = cxtTagsConditions; | ||
this.taskTypeCondition = taskTypeCondition; | ||
this.dataSourceCondition = dataSourceCondition; | ||
} | ||
|
||
/** | ||
* Evaluates this selector against a given task. | ||
* | ||
* @param task the task to evaluate | ||
* @return true if the task meets all the conditions specified by this selector, otherwise false | ||
*/ | ||
public boolean evaluate(Task task) | ||
{ | ||
boolean isMatch = true; | ||
if (cxtTagsConditions != null) { | ||
isMatch = cxtTagsConditions.entrySet().stream().allMatch(entry -> { | ||
String tagKey = entry.getKey(); | ||
Set<String> tagValues = entry.getValue(); | ||
Map<String, Object> tags = task.getContextValue(DruidMetrics.TAGS); | ||
if (tags == null || tags.isEmpty()) { | ||
return false; | ||
} | ||
Object tagValue = tags.get(tagKey); | ||
|
||
return tagValue == null ? false : tagValues.contains((String) tagValue); | ||
}); | ||
} | ||
|
||
if (isMatch && taskTypeCondition != null) { | ||
isMatch = taskTypeCondition.contains(task.getType()); | ||
} | ||
|
||
if (isMatch && dataSourceCondition != null) { | ||
isMatch = dataSourceCondition.contains(task.getDataSource()); | ||
} | ||
|
||
return isMatch; | ||
} | ||
|
||
@JsonProperty | ||
public String getSelectionKey() | ||
{ | ||
return selectionKey; | ||
} | ||
|
||
@JsonProperty("context.tags") | ||
public Map<String, Set<String>> getCxtTagsConditions() | ||
{ | ||
return cxtTagsConditions; | ||
} | ||
|
||
@JsonProperty("type") | ||
public Set<String> getTaskTypeCondition() | ||
{ | ||
return taskTypeCondition; | ||
} | ||
|
||
@JsonProperty("dataSource") | ||
public Set<String> getDataSourceCondition() | ||
{ | ||
return dataSourceCondition; | ||
} | ||
|
||
@Override | ||
public boolean equals(Object o) | ||
{ | ||
if (this == o) { | ||
return true; | ||
} | ||
if (o == null || getClass() != o.getClass()) { | ||
return false; | ||
} | ||
Selector selector = (Selector) o; | ||
return Objects.equals(selectionKey, selector.selectionKey) && Objects.equals( | ||
cxtTagsConditions, | ||
selector.cxtTagsConditions | ||
) && Objects.equals(taskTypeCondition, selector.taskTypeCondition) && Objects.equals( | ||
dataSourceCondition, | ||
selector.dataSourceCondition | ||
); | ||
} | ||
|
||
@Override | ||
public int hashCode() | ||
{ | ||
return Objects.hash(selectionKey, cxtTagsConditions, taskTypeCondition, dataSourceCondition); | ||
} | ||
|
||
@Override | ||
public String toString() | ||
{ | ||
return "Selector{" + | ||
"selectionKey=" + selectionKey + | ||
", context.tags=" + cxtTagsConditions + | ||
", type=" + taskTypeCondition + | ||
", dataSource=" + dataSourceCondition + | ||
'}'; | ||
} | ||
} |
112 changes: 112 additions & 0 deletions
112
.../java/org/apache/druid/k8s/overlord/execution/SelectorBasedPodTemplateSelectStrategy.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,112 @@ | ||
/* | ||
* 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.druid.k8s.overlord.execution; | ||
|
||
import com.fasterxml.jackson.annotation.JsonCreator; | ||
import com.fasterxml.jackson.annotation.JsonProperty; | ||
import com.google.common.base.Preconditions; | ||
import io.fabric8.kubernetes.api.model.PodTemplate; | ||
import org.apache.druid.indexing.common.task.Task; | ||
|
||
import javax.annotation.Nullable; | ||
import java.util.List; | ||
import java.util.Map; | ||
import java.util.Objects; | ||
|
||
/** | ||
* Implements {@link PodTemplateSelectStrategy} by dynamically evaluating a series of selectors. | ||
* Each selector corresponds to a potential task template key. | ||
*/ | ||
public class SelectorBasedPodTemplateSelectStrategy implements PodTemplateSelectStrategy | ||
{ | ||
@Nullable | ||
private String defaultKey; | ||
private List<Selector> selectors; | ||
|
||
@JsonCreator | ||
public SelectorBasedPodTemplateSelectStrategy( | ||
@JsonProperty("selectors") List<Selector> selectors, | ||
@JsonProperty("defaultKey") @Nullable String defaultKey | ||
) | ||
{ | ||
Preconditions.checkNotNull(selectors, "selectors"); | ||
this.selectors = selectors; | ||
this.defaultKey = defaultKey; | ||
} | ||
|
||
/** | ||
* Evaluates the provided task against the set selectors to determine its template. | ||
* | ||
* @param task the task to be checked | ||
* @return the template if a selector matches, otherwise fallback to base template | ||
*/ | ||
@Override | ||
public PodTemplate getPodTemplateForTask(Task task, Map<String, PodTemplate> templates) | ||
{ | ||
String templateKey = selectors.stream() | ||
.filter(selector -> selector.evaluate(task)) | ||
.findFirst() | ||
.map(Selector::getSelectionKey) | ||
.orElse(defaultKey); | ||
|
||
return templates.getOrDefault(templateKey, templates.get("base")); | ||
} | ||
|
||
@JsonProperty | ||
public List<Selector> getSelectors() | ||
{ | ||
return selectors; | ||
} | ||
|
||
@Nullable | ||
@JsonProperty | ||
public String getDefaultKey() | ||
{ | ||
return defaultKey; | ||
} | ||
|
||
@Override | ||
public boolean equals(Object o) | ||
{ | ||
if (this == o) { | ||
return true; | ||
} | ||
if (o == null || getClass() != o.getClass()) { | ||
return false; | ||
} | ||
SelectorBasedPodTemplateSelectStrategy that = (SelectorBasedPodTemplateSelectStrategy) o; | ||
return Objects.equals(defaultKey, that.defaultKey) && Objects.equals(selectors, that.selectors); | ||
} | ||
|
||
@Override | ||
public int hashCode() | ||
{ | ||
return Objects.hash(defaultKey, selectors); | ||
} | ||
|
||
@Override | ||
public String toString() | ||
{ | ||
return "SelectorBasedPodTemplateSelectStrategy{" + | ||
"selectors=" + selectors + | ||
", defaultKey=" + defaultKey + | ||
'}'; | ||
} | ||
} |
Oops, something went wrong.