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)
}