wilfred-s commented on code in PR #1062:
URL: https://github.com/apache/yunikorn-k8shim/pull/1062#discussion_r3772855042
##########
pkg/cache/task.go:
##########
@@ -343,46 +353,68 @@ func (task *Task) postTaskPending() {
// This routine binds the pod to the allocated node.
// It calls K8s api to bind a pod to the assigned node, this may need some
time,
// so we do a delay binding, background process, to avoid blocking main
process.
-// The result of the binding is tracked and failures are properly handled.
-// If successful, we move task to next state BOUND, otherwise we fail the task
+// Volume binding and pod binding are retried with a backoff; if they
ultimately fail
+// the allocation is rolled back to a pending ask so the core can re-schedule
the task
+// on a different node. On success we move the task to the next state BOUND.
func (task *Task) postTaskAllocated() {
- go func() {
- // we need to obtain task's lock first,
- // this ensures no other threads modifying task state at the
time being
- task.lock.Lock()
- defer task.lock.Unlock()
+ // Snapshot the fields needed for binding before launching the
goroutine without re-acquiring the lock.
+ // We already hold task.lock (via task.handle()) during state
transitions.
+ pod := task.pod
+ alias := task.alias
+ nodeName := task.nodeName
+ allocationKey := task.allocationKey
+ go func(pod *v1.Pod, alias, nodeName, allocationKey string) {
// post a message to indicate the pod gets its allocation
- events.GetRecorder().Eventf(task.pod.DeepCopy(),
+ events.GetRecorder().Eventf(pod.DeepCopy(),
nil, v1.EventTypeNormal, "Scheduled", "Scheduled",
- "Successfully assigned %s to node %s", task.alias,
task.nodeName)
+ "Successfully assigned %s to node %s", alias, nodeName)
Review Comment:
This is not really successful as yet, the attempt is being made but that
might fail.
This should be moved to the point we have done the bind and we log it in
line L407
##########
pkg/cache/task.go:
##########
@@ -625,11 +657,20 @@ func (task *Task) failWithEvent(errorMessage,
actionReason string) {
// pending ask so it can be re-scheduled on a different node.
// Must be called without holding the task lock.
func (task *Task) rollbackOnAssumePodFailure(allocationKey, nodeID string) {
- // Read fields needed for event posting and release request.
+ task.lock.RLock()
+ alias := task.alias
+ task.lock.RUnlock()
+ task.rollbackAllocation(allocationKey, nodeID, "AssumePodFailed",
+ fmt.Sprintf("Node assignment failed for %s on node %s, it will
be retried", alias, nodeID))
+}
+
+// rollbackAllocation resets the task allocation state and notifies the core
to move
+// the allocation back to a pending ask so it can be re-scheduled on a
different node.
+// Must be called without holding the task lock.
+func (task *Task) rollbackAllocation(allocationKey, nodeID, eventReason,
eventMsg string) {
Review Comment:
This is now called from the locked go routine, no locking here.
Comment is correct, code is not.
##########
pkg/cache/task.go:
##########
@@ -625,11 +657,20 @@ func (task *Task) failWithEvent(errorMessage,
actionReason string) {
// pending ask so it can be re-scheduled on a different node.
// Must be called without holding the task lock.
func (task *Task) rollbackOnAssumePodFailure(allocationKey, nodeID string) {
- // Read fields needed for event posting and release request.
+ task.lock.RLock()
+ alias := task.alias
+ task.lock.RUnlock()
Review Comment:
Please check locking as this is called from an unlocked location. The callee
`rollbackAllocation` needs a write lock to be held.
##########
pkg/cache/task.go:
##########
@@ -343,46 +353,68 @@ func (task *Task) postTaskPending() {
// This routine binds the pod to the allocated node.
// It calls K8s api to bind a pod to the assigned node, this may need some
time,
// so we do a delay binding, background process, to avoid blocking main
process.
-// The result of the binding is tracked and failures are properly handled.
-// If successful, we move task to next state BOUND, otherwise we fail the task
+// Volume binding and pod binding are retried with a backoff; if they
ultimately fail
+// the allocation is rolled back to a pending ask so the core can re-schedule
the task
+// on a different node. On success we move the task to the next state BOUND.
func (task *Task) postTaskAllocated() {
- go func() {
- // we need to obtain task's lock first,
- // this ensures no other threads modifying task state at the
time being
- task.lock.Lock()
- defer task.lock.Unlock()
+ // Snapshot the fields needed for binding before launching the
goroutine without re-acquiring the lock.
+ // We already hold task.lock (via task.handle()) during state
transitions.
+ pod := task.pod
+ alias := task.alias
+ nodeName := task.nodeName
+ allocationKey := task.allocationKey
+ go func(pod *v1.Pod, alias, nodeName, allocationKey string) {
// post a message to indicate the pod gets its allocation
- events.GetRecorder().Eventf(task.pod.DeepCopy(),
+ events.GetRecorder().Eventf(pod.DeepCopy(),
nil, v1.EventTypeNormal, "Scheduled", "Scheduled",
- "Successfully assigned %s to node %s", task.alias,
task.nodeName)
+ "Successfully assigned %s to node %s", alias, nodeName)
// before binding pod to node, first bind volumes to pod
log.Log(log.ShimCacheTask).Debug("bind pod volumes",
- zap.String("podName", task.pod.Name),
- zap.String("podUID", string(task.pod.UID)))
- if err := task.context.bindPodVolumes(task.pod); err != nil {
- log.Log(log.ShimCacheTask).Error("bind volumes to pod
failed", zap.String("taskID", task.taskID), zap.Error(err))
- task.failWithEvent(fmt.Sprintf("bind volumes to pod
failed, name: %s, %s", task.alias, err.Error()), "PodVolumesBindFailure")
+ zap.String("podName", pod.Name),
+ zap.String("podUID", string(pod.UID)))
+ if err := retry.OnError(bindPodBackoff, func(err error) bool {
+ log.Log(log.ShimCacheTask).Error("bind volumes to pod
failed, retrying",
+ zap.String("taskID", task.taskID),
zap.Error(err))
+ return true
+ }, func() error {
+ return task.context.bindPodVolumes(pod)
+ }); err != nil {
+ log.Log(log.ShimCacheTask).Error("bind volumes to pod
failed after retries",
+ zap.String("taskID", task.taskID),
zap.Error(err))
+ task.rescheduleOnBindFailure(allocationKey, nodeName,
"PodVolumesBindFailure",
+ fmt.Sprintf("Failed to bind volumes for %s on
node %s, it will be retried", alias, nodeName))
return
}
log.Log(log.ShimCacheTask).Debug("bind pod",
- zap.String("podName", task.pod.Name),
- zap.String("podUID", string(task.pod.UID)))
+ zap.String("podName", pod.Name),
+ zap.String("podUID", string(pod.UID)))
- if err :=
task.context.apiProvider.GetAPIs().KubeClient.Bind(task.pod, task.nodeName);
err != nil {
- log.Log(log.ShimCacheTask).Error("bind pod to node
failed", zap.String("taskID", task.taskID), zap.Error(err))
- task.failWithEvent(fmt.Sprintf("bind pod to node
failed, name: %s, %s", task.alias, err.Error()), "PodBindFailure")
+ if err := retry.OnError(bindPodBackoff, func(err error) bool {
+ log.Log(log.ShimCacheTask).Error("bind pod to node
failed, retrying",
+ zap.String("taskID", task.taskID),
zap.Error(err))
+ return true
+ }, func() error {
+ return
task.context.apiProvider.GetAPIs().KubeClient.Bind(pod, nodeName)
+ }); err != nil {
+ log.Log(log.ShimCacheTask).Error("bind pod to node
failed after retries",
+ zap.String("taskID", task.taskID),
zap.Error(err))
+ task.rescheduleOnBindFailure(allocationKey, nodeName,
"PodBindFailure",
+ fmt.Sprintf("Failed to bind %s to node %s, it
will be retried", alias, nodeName))
return
}
- log.Log(log.ShimCacheTask).Info("successfully bound pod",
zap.String("podName", task.pod.Name))
- dispatcher.Dispatch(NewBindTaskEvent(task.applicationID,
task.taskID))
- events.GetRecorder().Eventf(task.pod.DeepCopy(), nil,
- v1.EventTypeNormal, "PodBindSuccessful",
"PodBindSuccessful",
- "Pod %s is successfully bound to node %s", task.alias,
task.nodeName)
+ log.Log(log.ShimCacheTask).Info("successfully bound pod",
zap.String("podName", pod.Name))
+ task.lock.Lock()
task.schedulingState = TaskSchedAllocated
- }()
+ task.lock.Unlock()
Review Comment:
This lock needs to be taken at the start of the go routine as we do not want
anything to manipulate the task while running. It has to be a write lock also.
The removal of the lock outside the routine is correct but we need the lock
inside the go routine as the first thing. The go routine will be blocked until
the `handle()` returns as expected
##########
pkg/cache/task.go:
##########
@@ -669,6 +709,22 @@ func (task *Task)
rollbackOnAssumePodFailure(allocationKey, nodeID string) {
zap.String("allocationKey", allocationKey))
}
+// rescheduleOnBindFailure is called when volume or pod binding fails after
all retries.
+// Move the task back to Scheduling before rolling back the allocation, so a
+// subsequent TaskAllocated event from the core (which requires Scheduling) is
accepted.
+// Must be called without holding the task lock.
+func (task *Task) rescheduleOnBindFailure(allocationKey, nodeID, eventReason,
eventMsg string) {
+ // Move the task back to Scheduling before releasing to the core, so
the re-delivered
+ // allocation (valid only from the Scheduling state) is accepted by the
state machine.
+ if err := task.handle(NewRescheduleTaskEvent(task.applicationID,
task.taskID)); err != nil {
Review Comment:
You must hold the task lock to do the (volume)binding. The calling go
routine has the write lock. The handle will block as it needs a read lock.
The pattern used in the k8shim is not good for this. You either need to
dispatch an event or simply change the state back via
`task.sm.SetState(states.Scheduling). Preference is a direct state setting as
it has less overhead than an event dispatch.
##########
pkg/cache/task_state.go:
##########
@@ -47,10 +47,11 @@ const (
TaskFail
KillTask
TaskKilled
+ TaskRescheduling
Review Comment:
See above, direct state setting is preferred.
--
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]