adityadtu5 commented on code in PR #1062:
URL: https://github.com/apache/yunikorn-k8shim/pull/1062#discussion_r3774227559


##########
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:
   agreed, making this change



-- 
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]

Reply via email to