Inside the Kubernetes Job Controller: Controller Setup and Event Handlers

Entry Point

The Job controller lives in pkg/controller/job/job_controller.go. NewController() builds the controller instance, and its Run() method starts the sync loops. Higher up, the controller manager wires it like this:

  • cmd/kube-controller-manager/app/batch.go:34
func startJobController(ctx context.Context, controllerContext ControllerContext) (controller.Interface, bool, error) {
    go job.NewController(
        controllerContext.InformerFactory.Core().V1().Pods(),
        controllerContext.InformerFactory.Batch().V1().Jobs(),
        controllerContext.ClientBuilder.ClientOrDie("job-controller"),
    ).Run(int(controllerContext.ComponentConfig.JobController.ConcurrentJobSyncs), ctx.Done())
    return nil, true, nil
}

It calls job.NewController() with the Pod Informer, Job Informer, and a Kubernetes client, then starts the main loop with the configured concurrency.

Creating the Controller

The Controller struct

func NewController(podInformer coreinformers.PodInformer, jobInformer batchinformers.JobInformer, kubeClient clientset.Interface) *Controller

The three arguments are the Pod informer, the Job informer, and a clientset. Internally they are used to register event handlers and to set up listers. The returned *Controller is defined as:

  • pkg/controller/job/job_controller.go:80
type Controller struct {
    kubeClient          clientset.Interface
    podCtrl             controller.PodControlInterface
    updateStatusHandler func(job *batch.Job) error
    patchJobHandler     func(job *batch.Job, patch []byte) error
    syncHandler         func(jobKey string) (bool, error)
    podStoreSynced      cache.InformerSynced
    jobStoreSynced      cache.InformerSynced
    expectations        controller.ControllerExpectationsInterface
    jobLister           batchv1listers.JobLister
    podStore            corelisters.PodLister
    queue               workqueue.RateLimitingInterface
    orphanQueue         workqueue.RateLimitingInterface
    recorder            record.EventRecorder
}

We’ll see each field in action shortly.

NewController() walkthrough

func NewController(podInformer coreinformers.PodInformer, jobInformer batchinformers.JobInformer, kubeClient clientset.Interface) *Controller {
    eventBroadcaster := record.NewBroadcaster()
    eventBroadcaster.StartStructuredLogging(0)
    eventBroadcaster.StartRecordingToSink(&v1core.EventSinkImpl{Interface: kubeClient.CoreV1().Events("")})

    ctrl := &Controller{
        kubeClient: kubeClient,
        podCtrl: controller.RealPodControl{
            KubeClient: kubeClient,
            Recorder:   eventBroadcaster.NewRecorder(scheme.Scheme, v1.EventSource{Component: "job-controller"}),
        },
        expectations: controller.NewControllerExpectations(),
        queue:        workqueue.NewNamedRateLimitingQueue(workqueue.NewItemExponentialFailureRateLimiter(DefaultJobBackOff, MaxJobBackOff), "job"),
        orphanQueue:  workqueue.NewNamedRateLimitingQueue(workqueue.NewItemExponentialFailureRateLimiter(DefaultJobBackOff, MaxJobBackOff), "job_orphan_pod"),
        recorder:     eventBroadcaster.NewRecorder(scheme.Scheme, v1.EventSource{Component: "job-controller"}),
    }

    jobInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc: func(obj interface{}) {
            ctrl.enqueueForSync(obj, true)
        },
        UpdateFunc: ctrl.updateJob,
        DeleteFunc: func(obj interface{}) {
            ctrl.enqueueForSync(obj, true)
        },
    })
    ctrl.jobLister = jobInformer.Lister()
    ctrl.jobStoreSynced = jobInformer.Informer().HasSynced

    podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc:    ctrl.addPod,
        UpdateFunc: ctrl.updatePod,
        DeleteFunc: ctrl.deletePod,
    })
    ctrl.podStore = podInformer.Lister()
    ctrl.podStoreSynced = podInformer.Informer().HasSynced

    ctrl.updateStatusHandler = ctrl.updateJobStatus
    ctrl.patchJobHandler = ctrl.patchJob
    ctrl.syncHandler = ctrl.syncJob

    metrics.Register()

    return ctrl
}

The podCtrl field (previously named podControl) implements controller.PodControlInterface:

type PodControlInterface interface {
    CreatePods(namespace string, template *v1.PodTemplateSpec, object runtime.Object, controllerRef *metav1.OwnerReference) error
    CreatePodsWithGenerateName(namespace string, template *v1.PodTemplateSpec, object runtime.Object, controllerRef *metav1.OwnerReference, generateName string) error
    DeletePod(namespace string, podID string, object runtime.Object) error
    PatchPod(namespace, name string, data []byte) error
}

The concrete implementation controller.RealPodControl uses the clientset and event recorder to actually create, delete, and patch Pods.

Event Handlers

Job Add and Delete handlers

Both AddFunc and DeleteFunc for Jobs simply call ctrl.enqueueForSync(obj, true). Here is the renamed helper:

  • pkg/controller/job/job_controller.go:417
func (ctrl *Controller) enqueueForSync(obj interface{}, force bool) {
    key, err := controller.KeyFunc(obj)
    if err != nil {
        utilruntime.HandleError(fmt.Errorf("could not extract key for object %+v: %v", obj, err))
        return
    }
    delay := time.Duration(0)
    if !force {
        delay = computeBackoff(ctrl.queue, key)
    }
    klog.Infof("enqueuing job %s", key)
    ctrl.queue.AddAfter(key, delay)
}

If force is false the item waits for a backoff duration calculated from the number of previosu requeues (10s * 2^(n-1)). For add and delete events force is true, so the item is added immediately.

Job Update handler

Updates do the same enqueue, but additionally handle ActiveDeadlineSeconds changes:

func (ctrl *Controller) updateJob(old, cur interface{}) {
    oldJob := old.(*batch.Job)
    curJob := cur.(*batch.Job)
    key, err := controller.KeyFunc(curJob)
    if err != nil {
        return
    }
    ctrl.enqueueForSync(curJob, true)

    if curJob.Status.StartTime != nil {
        curADS := curJob.Spec.ActiveDeadlineSeconds
        if curADS == nil {
            return
        }
        oldADS := oldJob.Spec.ActiveDeadlineSeconds
        if oldADS == nil || *oldADS != *curADS {
            now := metav1.Now()
            start := curJob.Status.StartTime.Time
            elapsed := now.Time.Sub(start)
            total := time.Duration(*curADS) * time.Second
            ctrl.queue.AddAfter(key, total-elapsed)
            klog.V(4).Infof("job %q ActiveDeadlineSeconds updated, will resync after %d seconds", key, total-elapsed)
        }
    }
}

This calculates the reamining time until the deadline and schedules a requeue so that the sync loop can react when the deadline expires.

Pod Add handler

  • pkg/controller/job/job_controller.go:232
func (ctrl *Controller) addPod(obj interface{}) {
    pod := obj.(*v1.Pod)
    if pod.DeletionTimestamp != nil {
        // on controller restart we may get an add for a deleting Pod
        ctrl.deletePod(pod)
        return
    }

    if controllerRef := metav1.GetControllerOf(pod); controllerRef != nil {
        job := ctrl.resolveControllerRef(pod.Namespace, controllerRef)
        if job == nil {
            return
        }
        jobKey, err := controller.KeyFunc(job)
        if err != nil {
            return
        }
        ctrl.expectations.CreationObserved(jobKey)
        ctrl.enqueueForSync(job, true)
        return
    }

    // orphan pod – notify candidate owners
    for _, job := range ctrl.getPodJobs(pod) {
        ctrl.enqueueForSync(job, true)
    }
}

Pod Update handler

func (ctrl *Controller) updatePod(old, cur interface{}) {
    curPod := cur.(*v1.Pod)
    oldPod := old.(*v1.Pod)
    if curPod.ResourceVersion == oldPod.ResourceVersion {
        return // periodic resync, skip
    }
    if curPod.DeletionTimestamp != nil {
        ctrl.deletePod(curPod)
        return
    }

    immediate := curPod.Status.Phase != v1.PodFailed
    finalizerRemoved := !hasJobTrackingFinalizer(curPod)
    curControllerRef := metav1.GetControllerOf(curPod)
    oldControllerRef := metav1.GetControllerOf(oldPod)
    controllerRefChanged := !reflect.DeepEqual(curControllerRef, oldControllerRef)

    if controllerRefChanged && oldControllerRef != nil {
        if job := ctrl.resolveControllerRef(oldPod.Namespace, oldControllerRef); job != nil {
            if finalizerRemoved {
                key, err := controller.KeyFunc(job)
                if err == nil {
                    ctrl.finalizerExpectations.finalizerRemovalObserved(key, string(curPod.UID))
                }
            }
            ctrl.enqueueForSync(job, immediate)
        }
    }

    if curControllerRef != nil {
        job := ctrl.resolveControllerRef(curPod.Namespace, curControllerRef)
        if job == nil {
            return
        }
        if finalizerRemoved {
            key, err := controller.KeyFunc(job)
            if err == nil {
                ctrl.finalizerExpectations.finalizerRemovalObserved(key, string(curPod.UID))
            }
        }
        ctrl.enqueueForSync(job, immediate)
        return
    }

    labelChanged := !reflect.DeepEqual(curPod.Labels, oldPod.Labels)
    if labelChanged || controllerRefChanged {
        for _, job := range ctrl.getPodJobs(curPod) {
            ctrl.enqueueForSync(job, immediate)
        }
    }
}

Pod Delete handler

func (ctrl *Controller) deletePod(obj interface{}) {
    pod, ok := obj.(*v1.Pod)
    if !ok {
        tombstone, ok := obj.(cache.DeletedFinalStateUnknown)
        if !ok {
            utilruntime.HandleError(fmt.Errorf("could not get object from tombstone %+v", obj))
            return
        }
        pod, ok = tombstone.Obj.(*v1.Pod)
        if !ok {
            utilruntime.HandleError(fmt.Errorf("tombstone contained object that is not a pod %+v", obj))
            return
        }
    }
    controllerRef := metav1.GetControllerOf(pod)
    if controllerRef == nil {
        return
    }
    job := ctrl.resolveControllerRef(pod.Namespace, controllerRef)
    if job == nil {
        if hasJobTrackingFinalizer(pod) {
            ctrl.enqueueOrphanPod(pod)
        }
        return
    }
    jobKey, err := controller.KeyFunc(job)
    if err != nil {
        return
    }
    ctrl.expectations.DeletionObserved(jobKey)
    ctrl.finalizerExpectations.finalizerRemovalObserved(jobKey, string(pod.UID))
    ctrl.enqueueForSync(job, true)
}

Tags: kubernetes job-controller Go informers Workqueue

Posted on Fri, 24 Jul 2026 16:34:42 +0000 by anolan13