wilfred-s commented on code in PR #1098:
URL: https://github.com/apache/yunikorn-core/pull/1098#discussion_r3502974523
##########
pkg/scheduler/partition.go:
##########
@@ -42,9 +42,12 @@ import (
"github.com/apache/yunikorn-core/pkg/scheduler/policies"
"github.com/apache/yunikorn-core/pkg/scheduler/ugm"
"github.com/apache/yunikorn-core/pkg/webservice/dao"
+ siCommon "github.com/apache/yunikorn-scheduler-interface/lib/go/common"
"github.com/apache/yunikorn-scheduler-interface/lib/go/si"
)
+const allocationRollbackTag = siCommon.DomainYuniKornInternal +
"allocation-rollback"
Review Comment:
Move to `common.constant.go`
##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc
*objects.Allocation) (request
return false, true, nil
}
+ if existing.IsAllocated() && !allocated &&
isAllocationRollbackRequest(alloc) {
+ log.Log(log.SchedPartition).Info("handling allocation rollback",
+ zap.String("partitionName", pc.Name),
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.String("nodeID", existing.GetNodeID()))
+
+ if existingNode.RemoveAllocation(allocationKey) == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from node %s", allocationKey, existing.GetNodeID())
+ }
+ // this only removes allocation from the allocations map and
the allocation is still present in requests map
+ removed := app.RemoveAllocation(allocationKey,
si.TerminationType_STOPPED_BY_RM)
+ if removed == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from application %s", allocationKey, applicationID)
+ }
+
+ if err :=
queue.DecAllocatedResource(existing.GetAllocatedResource()); err != nil {
+ log.Log(log.SchedPartition).Warn("failed to release
resources from queue during rollback",
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.Error(err))
+ }
+
metrics.GetQueueMetrics(queue.GetQueuePath()).IncReleasedContainer()
+
+ if _, err := app.DeallocateAsk(allocationKey); err != nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to deallocate
ask %s on app %s during rollback: %w", allocationKey, applicationID, err)
+ }
+
+ if app.IsCompleting() {
+ if err :=
app.HandleApplicationEvent(objects.RunApplication); err != nil {
+ log.Log(log.SchedPartition).Warn("failed to
transition application back to running after rollback",
+ zap.String("appID", applicationID),
+ zap.String("allocationKey",
allocationKey),
+ zap.Error(err))
+ }
+ }
Review Comment:
This failure should never trigger a state change. The change is state
agnostic: the allocation is still there now in pending form.
##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc
*objects.Allocation) (request
return false, true, nil
}
+ if existing.IsAllocated() && !allocated &&
isAllocationRollbackRequest(alloc) {
+ log.Log(log.SchedPartition).Info("handling allocation rollback",
+ zap.String("partitionName", pc.Name),
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.String("nodeID", existing.GetNodeID()))
+
+ if existingNode.RemoveAllocation(allocationKey) == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from node %s", allocationKey, existing.GetNodeID())
+ }
+ // this only removes allocation from the allocations map and
the allocation is still present in requests map
+ removed := app.RemoveAllocation(allocationKey,
si.TerminationType_STOPPED_BY_RM)
+ if removed == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from application %s", allocationKey, applicationID)
+ }
+
+ if err :=
queue.DecAllocatedResource(existing.GetAllocatedResource()); err != nil {
+ log.Log(log.SchedPartition).Warn("failed to release
resources from queue during rollback",
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.Error(err))
+ }
+
metrics.GetQueueMetrics(queue.GetQueuePath()).IncReleasedContainer()
+
+ if _, err := app.DeallocateAsk(allocationKey); err != nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to deallocate
ask %s on app %s during rollback: %w", allocationKey, applicationID, err)
+ }
+
+ if app.IsCompleting() {
+ if err :=
app.HandleApplicationEvent(objects.RunApplication); err != nil {
+ log.Log(log.SchedPartition).Warn("failed to
transition application back to running after rollback",
+ zap.String("appID", applicationID),
+ zap.String("allocationKey",
allocationKey),
+ zap.Error(err))
+ }
+ }
+
+ existing.SetNodeID("")
+ existing.SetBindTime(time.Time{})
+ pc.updateAllocationCount(-1)
+ if existing.IsPlaceholder() {
+ pc.decPhAllocationCount(1)
+ }
Review Comment:
Placeholders cannot specify volumes. They cannot fail due to volume bind
failures.
##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc
*objects.Allocation) (request
return false, true, nil
}
+ if existing.IsAllocated() && !allocated &&
isAllocationRollbackRequest(alloc) {
+ log.Log(log.SchedPartition).Info("handling allocation rollback",
+ zap.String("partitionName", pc.Name),
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.String("nodeID", existing.GetNodeID()))
+
+ if existingNode.RemoveAllocation(allocationKey) == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from node %s", allocationKey, existing.GetNodeID())
+ }
+ // this only removes allocation from the allocations map and
the allocation is still present in requests map
+ removed := app.RemoveAllocation(allocationKey,
si.TerminationType_STOPPED_BY_RM)
Review Comment:
This is not correct: we do have two maps with all allocated and pending
allocations but they are the same object. When you remove it here it might be
gone from the map but it is not released for scheduling.
We need a specific rollback for the app and queue:
* decrease allocated resource on app, user and queue
* increase pending resources on app and queue
* remove the reference from the allocated map
* remove the node from the allocation, unset allocated flag (must become
false)
* add an entry to the allocLog of the allocation for the bind failure
Some of these steps are scattered around in this code but it is not
correct/complete
##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc
*objects.Allocation) (request
return false, true, nil
}
+ if existing.IsAllocated() && !allocated &&
isAllocationRollbackRequest(alloc) {
+ log.Log(log.SchedPartition).Info("handling allocation rollback",
+ zap.String("partitionName", pc.Name),
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.String("nodeID", existing.GetNodeID()))
+
+ if existingNode.RemoveAllocation(allocationKey) == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from node %s", allocationKey, existing.GetNodeID())
+ }
+ // this only removes allocation from the allocations map and
the allocation is still present in requests map
+ removed := app.RemoveAllocation(allocationKey,
si.TerminationType_STOPPED_BY_RM)
+ if removed == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from application %s", allocationKey, applicationID)
+ }
+
+ if err :=
queue.DecAllocatedResource(existing.GetAllocatedResource()); err != nil {
+ log.Log(log.SchedPartition).Warn("failed to release
resources from queue during rollback",
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.Error(err))
+ }
+
metrics.GetQueueMetrics(queue.GetQueuePath()).IncReleasedContainer()
Review Comment:
The container is not released it never ran, this metric should not increase
##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc
*objects.Allocation) (request
return false, true, nil
}
+ if existing.IsAllocated() && !allocated &&
isAllocationRollbackRequest(alloc) {
+ log.Log(log.SchedPartition).Info("handling allocation rollback",
+ zap.String("partitionName", pc.Name),
+ zap.String("appID", applicationID),
+ zap.String("allocationKey", allocationKey),
+ zap.String("nodeID", existing.GetNodeID()))
+
+ if existingNode.RemoveAllocation(allocationKey) == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from node %s", allocationKey, existing.GetNodeID())
+ }
+ // this only removes allocation from the allocations map and
the allocation is still present in requests map
+ removed := app.RemoveAllocation(allocationKey,
si.TerminationType_STOPPED_BY_RM)
+ if removed == nil {
+ metrics.GetSchedulerMetrics().IncSchedulingError()
+ return false, false, fmt.Errorf("failed to remove
allocation %s from application %s", allocationKey, applicationID)
+ }
Review Comment:
This cannot happen as `existing` cannot be nil when we get here
--
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]