wilfred-s commented on code in PR #1106:
URL: https://github.com/apache/yunikorn-core/pull/1106#discussion_r4152550136
##########
pkg/scheduler/objects/predicates.go:
##########
@@ -113,9 +112,7 @@ func (pcr *predicateCheckResult)
populateVictims(victimsByNode map[string][]*All
}
}
-// preemptPredicateCheck performs a single predicate check and reports the
resultType on a channel
-func preemptPredicateCheck(plugin api.ResourceManagerCallback, ch chan<-
*predicateCheckResult, wg *sync.WaitGroup, args *si.PreemptionPredicatesArgs) {
- defer wg.Done()
+func PredicateChecks(plugin api.ResourceManagerCallback, args
*si.PreemptionPredicatesArgs) *predicateCheckResult {
Review Comment:
Internal to objects package no need to export (same as @stantheman0128
mentioned)
##########
pkg/scheduler/objects/required_node_preemptor.go:
##########
@@ -67,39 +68,41 @@ func (p *PreemptionContext) tryPreemption() {
p.sortAllocations()
// Are there any victims/asks to preempt?
- victims := p.GetVictims()
+ var victims []*Allocation
+ var finalVictims []*Allocation
+ victims = p.GetVictims()
if len(victims) > 0 {
- log.Log(log.SchedRequiredNodePreemption).Info("Found victims
for required node preemption",
- zap.String("ds allocation key",
p.requiredAsk.GetAllocationKey()),
- zap.String("allocation name",
p.requiredAsk.GetAllocationName()),
- zap.Int("no.of victims", len(victims)))
- for _, victim := range victims {
- err := victim.MarkPreempted()
- if err != nil {
-
log.Log(log.SchedRequiredNodePreemption).Warn("allocation is already released,
so not proceeding further on the daemon set preemption process",
- zap.String("applicationID",
p.requiredAsk.GetApplicationID()),
- zap.String("allocationKey",
victim.GetAllocationKey()))
- continue
+ finalVictims = p.runPredicates(victims)
+ if len(finalVictims) > 0 {
+ for _, victim := range finalVictims {
+ err := victim.MarkPreempted()
+ if err != nil {
+
log.Log(log.SchedRequiredNodePreemption).Warn("allocation is already released,
so not proceeding further on the daemon set preemption process",
+ zap.String("applicationID",
p.requiredAsk.GetApplicationID()),
+ zap.String("allocationKey",
victim.GetAllocationKey()))
+ continue
+ }
+ if victimQueue :=
p.application.queue.GetQueueByAppID(victim.GetApplicationID()); victimQueue !=
nil {
+
victimQueue.IncPreemptingResource(victim.GetAllocatedResource())
+ } else {
+
log.Log(log.SchedRequiredNodePreemption).Warn("BUG: Queue not found for daemon
set preemption victim",
+ zap.String("queue",
p.application.queue.Name),
+
zap.String("victimApplicationID", victim.GetApplicationID()),
+
zap.String("victimAllocationKey", victim.GetAllocationKey()))
+ }
+
victim.SendPreemptedBySchedulerEvent(p.requiredAsk.GetAllocationKey(),
p.requiredAsk.GetApplicationID(), p.application.queuePath)
}
- if victimQueue :=
p.application.queue.GetQueueByAppID(victim.GetApplicationID()); victimQueue !=
nil {
-
victimQueue.IncPreemptingResource(victim.GetAllocatedResource())
- } else {
-
log.Log(log.SchedRequiredNodePreemption).Warn("BUG: Queue not found for daemon
set preemption victim",
- zap.String("queue",
p.application.queue.Name),
- zap.String("victimApplicationID",
victim.GetApplicationID()),
- zap.String("victimAllocationKey",
victim.GetAllocationKey()))
- }
-
victim.SendPreemptedBySchedulerEvent(p.requiredAsk.GetAllocationKey(),
p.requiredAsk.GetApplicationID(), p.application.queuePath)
+ p.requiredAsk.MarkTriggeredPreemption()
+ p.application.notifyRMAllocationReleased(victims,
si.TerminationType_PREEMPTED_BY_SCHEDULER,
+ "preempting allocations to free up resources to
run daemon set ask: "+p.requiredAsk.GetAllocationKey())
}
- p.requiredAsk.MarkTriggeredPreemption()
- p.application.notifyRMAllocationReleased(victims,
si.TerminationType_PREEMPTED_BY_SCHEDULER,
- "preempting allocations to free up resources to run
daemon set ask: "+p.requiredAsk.GetAllocationKey())
- } else {
+ }
+ if len(victims) == 0 || len(finalVictims) == 0 {
Review Comment:
This can (or should) be just the check for `finalVictims == 0` the number of
victims that was originally found is no longer relevant after the predicates
removed discarded them.
If original victims length is 0 the final length is also 0
##########
pkg/scheduler/objects/required_node_preemptor.go:
##########
@@ -171,14 +174,68 @@ func (p *PreemptionContext) GetVictims() []*Allocation {
}
}
- // Did we found the useful set of victims?
+ // Did we find the useful set of victims?
if len(victims) > 0 && resources.StrictlyGreaterThanOrEquals(
Review Comment:
this needs to be the final list from after the predicate run
##########
pkg/scheduler/objects/required_node_preemptor.go:
##########
@@ -171,14 +174,68 @@ func (p *PreemptionContext) GetVictims() []*Allocation {
}
}
- // Did we found the useful set of victims?
+ // Did we find the useful set of victims?
if len(victims) > 0 && resources.StrictlyGreaterThanOrEquals(
resources.Add(currentResource, p.node.GetAvailableResource()),
p.requiredAsk.GetAllocatedResource()) {
return victims
}
return nil
}
+// runPredicates Run Predicate checks to confirm whether collected victims
from the required node is
+// good enough to move forward on the preemption further or not
+func (p *PreemptionContext) runPredicates(victims []*Allocation) []*Allocation
{
+ finalVictims := make([]*Allocation, 0)
+ plugin := plugins.GetResourceManagerCallbackPlugin()
+ if plugin == nil {
+ log.Log(log.SchedRequiredNodePreemption).Debug("No RM callback
plugin registered, using chosen victims as is",
+ zap.String("node", p.node.NodeID),
+ zap.String("allocationKey",
p.requiredAsk.GetAllocationKey()))
+ return victims
+ } else {
+ // run predicates for this pod before in hand and fetch
feasible nodes
+ feasibleNodes, predicatesResult :=
p.requiredAsk.preAllocateConditions(true)
Review Comment:
We need a rebase as this was added in #1096 and this PR was created before
that commit.
##########
pkg/scheduler/objects/required_node_preemptor.go:
##########
@@ -171,14 +174,68 @@ func (p *PreemptionContext) GetVictims() []*Allocation {
}
}
- // Did we found the useful set of victims?
+ // Did we find the useful set of victims?
if len(victims) > 0 && resources.StrictlyGreaterThanOrEquals(
resources.Add(currentResource, p.node.GetAvailableResource()),
p.requiredAsk.GetAllocatedResource()) {
return victims
Review Comment:
same here finalVictims
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]