wilfred-s commented on code in PR #1104:
URL: https://github.com/apache/yunikorn-core/pull/1104#discussion_r3626962525
##########
pkg/scheduler/partition.go:
##########
@@ -1504,6 +1504,45 @@ func (pc *PartitionContext) removeAllocation(release
*si.AllocationRelease) ([]*
return nil, nil
}
+ // Handle shim-initiated scheduling failure: roll the allocation back
to a pending ask
Review Comment:
Can we factor out this code in its own method?
##########
pkg/scheduler/objects/application.go:
##########
@@ -799,6 +810,51 @@ func (sa *Application) DeallocateAsk(allocKey string)
(*resources.Resource, erro
return nil, fmt.Errorf("failed to locate ask with key %s", allocKey)
}
+// RollbackAllocation atomically reverts an allocated ask back to pending
state.
+// This is used when the shim fails to bind a pod after the core has allocated
it,
+// allowing the task to be re-scheduled on a different node.
+// The application must be in Accepted or Running state for rollback to
succeed.
+func (sa *Application) RollbackAllocation(allocKey string)
(*resources.Resource, error) {
+ sa.Lock()
+ defer sa.Unlock()
+
+ if !sa.IsAccepted() && !sa.IsRunning() {
+ return nil, fmt.Errorf("cannot rollback allocation %s:
application %s is in state %s", allocKey, sa.ApplicationID, sa.CurrentState())
+ }
+
+ ask := sa.requests[allocKey]
+ if ask == nil {
+ return nil, fmt.Errorf("failed to locate ask with key %s for
rollback", allocKey)
+ }
+ alloc := sa.allocations[allocKey]
+ if alloc == nil {
+ return nil, fmt.Errorf("failed to locate allocation with key %s
for rollback", allocKey)
+ }
Review Comment:
The same reference (pointer) is used in both maps as the value.
This should always be true: `ask == alloc`
allocations must always be a subset of requests checking allocations should
be enough
##########
pkg/scheduler/partition.go:
##########
@@ -1504,6 +1504,45 @@ func (pc *PartitionContext) removeAllocation(release
*si.AllocationRelease) ([]*
return nil, nil
}
+ // Handle shim-initiated scheduling failure: roll the allocation back
to a pending ask
+ // so the core can re-schedule it on a different node. This bypasses
the normal
+ // remove-and-destroy path entirely. We return nil,nil to suppress the
echo back to
+ // the shim (the shim initiated this release; no confirmation is
needed).
+ if release.TerminationType ==
si.TerminationType_SCHEDULING_FAILED_ON_RM {
+ queue := app.GetQueue()
+ // Retrieve node ID before rolling back (RollbackAllocation
clears it on the ask).
+ nodeID := app.GetAllocationNodeID(allocationKey)
+ res, err := app.RollbackAllocation(allocationKey)
+ if err != nil {
+ log.Log(log.SchedPartition).Warn("failed to rollback
allocation",
+ zap.String("appID", appID),
+ zap.String("allocationKey", allocationKey),
+ zap.Error(err))
+ return nil, nil
+ }
+ if node := pc.GetNode(nodeID); node != nil {
+ node.RemoveAllocation(allocationKey)
+ } else {
+ log.Log(log.SchedPartition).Warn("node not found while
rolling back allocation",
+ zap.String("appID", appID),
+ zap.String("allocationKey", allocationKey),
+ zap.String("nodeID", nodeID))
+ }
+ if err := queue.DecAllocatedResource(res); err != nil {
+ log.Log(log.SchedPartition).Warn("failed to release
resources from queue during rollback",
+ zap.String("appID", appID),
+ zap.String("allocationKey", allocationKey),
+ zap.Error(err))
+ }
+ pc.updateAllocationCount(-1)
+
metrics.GetQueueMetrics(queue.GetQueuePath()).AddReleasedContainers(1)
Review Comment:
Two points:
* are we sure that the container was added to the allocated containers
metric already? The metric gets increased when `UpdateAllocation()` is called
after confirming the allocation. I think that has not happened yet.
* if above is not true: why not use `IncReleasedContainer()` ? less overhead
--
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]