wilfred-s commented on code in PR #1062:
URL: https://github.com/apache/yunikorn-k8shim/pull/1062#discussion_r3794009261
##########
pkg/cache/task.go:
##########
@@ -343,46 +354,73 @@ 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()
-
- // post a message to indicate the pod gets its allocation
- events.GetRecorder().Eventf(task.pod.DeepCopy(),
- nil, v1.EventTypeNormal, "Scheduled", "Scheduled",
- "Successfully assigned %s to node %s", task.alias,
task.nodeName)
+ // 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) {
+ // this lock is to protect the task from being modified while
we are binding the pod to node
+ // once all task related operations are done, release the lock
+ // This is important specially while calling
rollbackAllocation, which needs to acquire the context lock, so we cannot hold
the task lock while calling it.
+ task.lock.Lock()
// 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(retryBackoff, 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)
Review Comment:
This has a timeout already which is 10 minutes. If we retry for that kind of
failure the process would take 8*1m + 2m for the retry here. That is too much
we need to retry API calls not the wait for the provisioning.
Worried about the backoff/retry on top of the long timeout.
--
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]