client-go如何实现Job删除完成后再执行新Job创建操作
Go语言基于client-go实现等待Batch Job完全删除的方案
实现原理
Kubernetes的资源删除为异步逻辑:调用Delete接口仅会给资源打上删除标记,实际的资源清理由对应控制器异步完成,因此直接调用Delete后立即创建同名Job会出现冲突。我们需要等待Job资源从集群中完全移除后再执行创建操作,可通过轮询检查或Watch监听两种方式实现。
方案1:轮询检查实现(推荐,兼容性好、逻辑简单)
import ( "context" "time" batchv1 "k8s.io/api/batch/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/api/errors" ) // DeleteBatchJobAndWait 删除指定Job并等待完全删除成功 func (k K8sClient) DeleteBatchJobAndWait(name string, namespace string, timeout time.Duration) error { // 先触发删除 err := k.K8sCS.BatchV1().Jobs(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{}) if err != nil { // 若Job本身不存在直接返回成功 if errors.IsNotFound(err) { return nil } return err } // 轮询检查配置:每2秒查一次,直到超时 checkInterval := 2 * time.Second timeoutCtx, cancel := context.WithTimeout(context.TODO(), timeout) defer cancel() ticker := time.NewTicker(checkInterval) defer ticker.Stop() for { select { case <-timeoutCtx.Done(): return timeoutCtx.Err() case <-ticker.C: _, err := k.K8sCS.BatchV1().Jobs(namespace).Get(context.TODO(), name, metav1.GetOptions{}) if err != nil { if errors.IsNotFound(err) { // Job已完全删除,返回成功 return nil } return err } } } }
使用示例
调用该方法并传入合理的超时时间,返回无错误后即可执行新Job创建逻辑:
err := k8sClient.DeleteBatchJobAndWait("my-job", "default", 30 * time.Second) if err != nil { // 处理删除失败/超时错误 log.Fatalf("删除Job失败: %v", err) } // 旧Job已完全清理,此处执行新Job创建逻辑
方案2:Watch监听实现(高效,延迟更低)
适合对等待延迟要求较高的场景,通过监听删除事件触发后续逻辑,无需轮询:
import ( "context" "fmt" "time" batchv1 "k8s.io/api/batch/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/watch" ) // DeleteBatchJobAndWatch 删除Job并通过Watch等待删除完成 func (k K8sClient) DeleteBatchJobAndWatch(name string, namespace string, timeout time.Duration) error { err := k.K8sCS.BatchV1().Jobs(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{}) if err != nil { if errors.IsNotFound(err) { return nil } return err } timeoutCtx, cancel := context.WithTimeout(context.TODO(), timeout) defer cancel() timeoutSec := int64(timeout.Seconds()) // 启动指定Job的监听 watcher, err := k.K8sCS.BatchV1().Jobs(namespace).Watch(timeoutCtx, metav1.ListOptions{ FieldSelector: fmt.Sprintf("metadata.name=%s", name), TimeoutSeconds: &timeoutSec, }) if err != nil { return err } defer watcher.Stop() // 处理监听事件 for event := range watcher.ResultChan() { switch event.Type { case watch.Deleted: // 收到删除事件,Job已清理完成 return nil case watch.Error: return fmt.Errorf("监听异常: %v", event.Object) } } return timeoutCtx.Err() }
可选配置:同步删除关联Pod
如果需要确保Job对应的Pod也完全清理,可在删除时设置级联删除策略为Foreground,修改Delete调用参数如下:
propagationPolicy := metav1.DeletePropagationForeground err := k.K8sCS.BatchV1().Jobs(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{ PropagationPolicy: &propagationPolicy, })
内容的提问来源于stack exchange,提问作者Abhishek Kumar
相关产品推荐
相关产品推荐

