Let executor type manages how to do cleanup for old kubeobjects (#1455)

Add CleanupOldExecutorObjects to executor type interface in order
to let an executor type manages how to clean up the resources it created.
This commit is contained in:
Ta-Ching Chen
2019-12-02 19:47:31 +08:00
committed by GitHub
parent 3f3b11ffbf
commit 275da18cf6
6 changed files with 107 additions and 75 deletions
+6 -3
View File
@@ -244,18 +244,22 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
executorTypes[gpm.GetTypeName()] = gpm executorTypes[gpm.GetTypeName()] = gpm
executorTypes[ndm.GetTypeName()] = ndm executorTypes[ndm.GetTypeName()] = ndm
if ok, _ := strconv.ParseBool(os.Getenv("ADOPT_EXISTING_RESOURCES")); ok { adoptExistingResources, _ := strconv.ParseBool(os.Getenv("ADOPT_EXISTING_RESOURCES"))
wg := &sync.WaitGroup{} wg := &sync.WaitGroup{}
for _, et := range executorTypes { for _, et := range executorTypes {
wg.Add(1) wg.Add(1)
go func(et executortype.ExecutorType) { go func(et executortype.ExecutorType) {
defer wg.Done() defer wg.Done()
if adoptExistingResources {
et.AdoptExistingResources() et.AdoptExistingResources()
}
et.CleanupOldExecutorObjects()
}(et) }(et)
} }
// set hard timeout for resource adoption // set hard timeout for resource adoption
// TODO: use context to control the waiting time once kubernetes client supports it.
util.WaitTimeout(wg, 30*time.Second) util.WaitTimeout(wg, 30*time.Second)
}
cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes) cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes)
@@ -264,7 +268,6 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
return err return err
} }
go reaper.CleanupOldExecutorObjects(logger, kubernetesClient, executorInstanceID)
go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30) go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30)
go api.Serve(port) go api.Serve(port)
go serveMetric(logger) go serveMetric(logger)
@@ -53,4 +53,7 @@ type ExecutorType interface {
// AdoptOrphanResources adopts existing resources created by the deleted executor. // AdoptOrphanResources adopts existing resources created by the deleted executor.
AdoptExistingResources() AdoptExistingResources()
// CleanupOldExecutorObjects cleans up resources created by old executor instances
CleanupOldExecutorObjects()
} }
@@ -100,6 +100,10 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
} }
depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(deployment) depl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Create(deployment)
if err != nil {
if k8s_err.IsAlreadyExists(err) {
depl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(deployName, metav1.GetOptions{})
}
if err != nil { if err != nil {
deploy.logger.Error("error while creating function deployment", deploy.logger.Error("error while creating function deployment",
zap.Error(err), zap.Error(err),
@@ -108,16 +112,13 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
zap.String("deployment_namespace", deployNamespace)) zap.String("deployment_namespace", deployNamespace))
return nil, err return nil, err
} }
}
if minScale > 0 { if minScale > 0 {
depl, err = deploy.waitForDeploy(depl, minScale, specializationTimeout) depl, err = deploy.waitForDeploy(depl, minScale, specializationTimeout)
} }
return depl, err return depl, err
} }
return nil, err return nil, err
} }
func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) error { func (deploy *NewDeploy) setupRBACObjs(deployNamespace string, fn *fv1.Function) error {
@@ -402,14 +403,17 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut
return existingHpa, err return existingHpa, err
} else if k8s_err.IsNotFound(err) { } else if k8s_err.IsNotFound(err) {
cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(hpa) cHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Create(hpa)
if err != nil {
if k8s_err.IsAlreadyExists(err) {
cHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(hpaName, metav1.GetOptions{})
}
if err != nil { if err != nil {
return nil, err return nil, err
} }
}
return cHpa, nil return cHpa, nil
} }
return nil, err return nil, err
} }
func (deploy *NewDeploy) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) { func (deploy *NewDeploy) getHpa(ns, name string) (*asv1.HorizontalPodAutoscaler, error) {
@@ -464,9 +468,14 @@ func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAn
return existingSvc, err return existingSvc, err
} else if k8s_err.IsNotFound(err) { } else if k8s_err.IsNotFound(err) {
svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(service) svc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Create(service)
if err != nil {
if k8s_err.IsAlreadyExists(err) {
svc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(svcName, metav1.GetOptions{})
}
if err != nil { if err != nil {
return nil, err return nil, err
} }
}
return svc, nil return svc, nil
} }
return nil, err return nil, err
@@ -43,6 +43,7 @@ import (
"github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/executor/executortype" "github.com/fission/fission/pkg/executor/executortype"
"github.com/fission/fission/pkg/executor/fscache" "github.com/fission/fission/pkg/executor/fscache"
"github.com/fission/fission/pkg/executor/reaper"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config" fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/throttler" "github.com/fission/fission/pkg/throttler"
"github.com/fission/fission/pkg/types" "github.com/fission/fission/pkg/types"
@@ -272,6 +273,35 @@ func (deploy *NewDeploy) AdoptExistingResources() {
wg.Wait() wg.Wait()
} }
func (deploy *NewDeploy) CleanupOldExecutorObjects() {
deploy.logger.Info("Newdeploy starts to clean orphaned resources", zap.String("instanceID", deploy.instanceID))
errs := &multierror.Error{}
listOpts := metav1.ListOptions{
LabelSelector: labels.Set(map[string]string{types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy)}).AsSelector().String(),
}
err := reaper.CleanupHpa(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupDeployments(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupServices(deploy.logger, deploy.kubernetesClient, deploy.instanceID, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
if errs.ErrorOrNil() != nil {
// TODO retry reaper; logged and ignored for now
deploy.logger.Error("Failed to cleanup old executor objects", zap.Error(err))
}
}
func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) { func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything()) listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything())
+25
View File
@@ -26,6 +26,7 @@ import (
"sync" "sync"
"time" "time"
"github.com/hashicorp/go-multierror"
"go.uber.org/zap" "go.uber.org/zap"
apiv1 "k8s.io/api/core/v1" apiv1 "k8s.io/api/core/v1"
k8serrors "k8s.io/apimachinery/pkg/api/errors" k8serrors "k8s.io/apimachinery/pkg/api/errors"
@@ -376,6 +377,30 @@ func (gpm *GenericPoolManager) AdoptExistingResources() {
wg.Wait() wg.Wait()
} }
func (gpm *GenericPoolManager) CleanupOldExecutorObjects() {
gpm.logger.Info("Poolmanager starts to clean orphaned resources", zap.String("instanceID", gpm.instanceId))
errs := &multierror.Error{}
listOpts := metav1.ListOptions{
LabelSelector: labels.Set(map[string]string{types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr)}).AsSelector().String(),
}
err := reaper.CleanupDeployments(gpm.logger, gpm.kubernetesClient, gpm.instanceId, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
err = reaper.CleanupPods(gpm.logger, gpm.kubernetesClient, gpm.instanceId, listOpts)
if err != nil {
errs = multierror.Append(errs, err)
}
if errs.ErrorOrNil() != nil {
// TODO retry reaper; logged and ignored for now
gpm.logger.Error("Failed to cleanup old executor objects", zap.Error(err))
}
}
func (gpm *GenericPoolManager) service() { func (gpm *GenericPoolManager) service() {
for { for {
req := <-gpm.requestChannel req := <-gpm.requestChannel
+12 -50
View File
@@ -20,7 +20,6 @@ import (
"strings" "strings"
"time" "time"
"github.com/hashicorp/go-multierror"
"go.uber.org/zap" "go.uber.org/zap"
apiv1 "k8s.io/api/core/v1" apiv1 "k8s.io/api/core/v1"
meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1" meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -37,43 +36,6 @@ var (
delOpt = meta_v1.DeleteOptions{PropagationPolicy: &deletePropagation} delOpt = meta_v1.DeleteOptions{PropagationPolicy: &deletePropagation}
) )
// CleanupOldExecutorObjects cleans up resources created by old executor instances
func CleanupOldExecutorObjects(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, instanceId string) {
errs := &multierror.Error{}
err := cleanupHpa(logger, kubernetesClient, instanceId)
if err != nil {
errs = multierror.Append(errs, err)
}
err = cleanupDeployments(logger, kubernetesClient, instanceId)
if err != nil {
errs = multierror.Append(errs, err)
}
// Pods might be running user functions, so we give them
// a few minutes before terminating them. This time is the
// maximum function runtime, plus the time a router might
// still route to an old instance, i.e. router cache expiry
// time.
time.Sleep(6 * time.Minute)
err = cleanupServices(logger, kubernetesClient, instanceId)
if err != nil {
errs = multierror.Append(errs, err)
}
err = cleanupPods(logger, kubernetesClient, instanceId)
if err != nil {
errs = multierror.Append(errs, err)
}
if errs.ErrorOrNil() != nil {
// TODO retry reaper; logged and ignored for now
logger.Error("Failed to cleanup old executor objects", zap.Error(err))
}
}
// CleanupKubeObject deletes given kubernetes object // CleanupKubeObject deletes given kubernetes object
func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) { func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kubeobj *apiv1.ObjectReference) {
switch strings.ToLower(kubeobj.Kind) { switch strings.ToLower(kubeobj.Kind) {
@@ -107,8 +69,8 @@ func CleanupKubeObject(logger *zap.Logger, kubeClient *kubernetes.Clientset, kub
} }
} }
func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { func CleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error {
deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) deploymentList, err := client.AppsV1().Deployments(meta_v1.NamespaceAll).List(listOps)
if err != nil { if err != nil {
return err return err
} }
@@ -119,7 +81,7 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan
id, ok = dep.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] id, ok = dep.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
} }
if ok && id != instanceId { if ok && id != instanceId {
logger.Debug("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name)) logger.Info("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name))
err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt) err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt)
if err != nil { if err != nil {
logger.Error("error cleaning up deployment", logger.Error("error cleaning up deployment",
@@ -133,8 +95,8 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan
return nil return nil
} }
func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { func CleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error {
podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) podList, err := client.CoreV1().Pods(meta_v1.NamespaceAll).List(listOps)
if err != nil { if err != nil {
return err return err
} }
@@ -145,7 +107,7 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st
id, ok = pod.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] id, ok = pod.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
} }
if ok && id != instanceId { if ok && id != instanceId {
logger.Debug("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name)) logger.Info("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name))
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil) err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil)
if err != nil { if err != nil {
logger.Error("error cleaning up pod", logger.Error("error cleaning up pod",
@@ -159,8 +121,8 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st
return nil return nil
} }
func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { func CleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error {
svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) svcList, err := client.CoreV1().Services(meta_v1.NamespaceAll).List(listOps)
if err != nil { if err != nil {
return err return err
} }
@@ -171,7 +133,7 @@ func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI
id, ok = svc.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] id, ok = svc.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
} }
if ok && id != instanceId { if ok && id != instanceId {
logger.Debug("cleaning up service", zap.String("service", svc.ObjectMeta.Name)) logger.Info("cleaning up service", zap.String("service", svc.ObjectMeta.Name))
err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(svc.ObjectMeta.Name, nil) err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(svc.ObjectMeta.Name, nil)
if err != nil { if err != nil {
logger.Error("error cleaning up service", logger.Error("error cleaning up service",
@@ -185,8 +147,8 @@ func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI
return nil return nil
} }
func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error { func CleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId string, listOps meta_v1.ListOptions) error {
hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(meta_v1.ListOptions{}) hpaList, err := client.AutoscalingV1().HorizontalPodAutoscalers(meta_v1.NamespaceAll).List(listOps)
if err != nil { if err != nil {
return err return err
} }
@@ -198,7 +160,7 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str
id, ok = hpa.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL] id, ok = hpa.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
} }
if ok && id != instanceId { if ok && id != instanceId {
logger.Debug("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name)) logger.Info("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name))
err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(hpa.ObjectMeta.Name, nil) err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(hpa.ObjectMeta.Name, nil)
if err != nil { if err != nil {
logger.Error("error cleaning up HPA", logger.Error("error cleaning up HPA",