Adopt existing orphan kubernetes resources when executor starts up (#1443)

Previously, once the executor is deleted for reasons (like upgrade or cluster scale-in),
the new executor deletes all existing resources created by the old executor and creates
new one. This mechanism becomes a problem when there are requests connecting to the
existing pods. Also in the worst case, the cluster may not have enough resources to create
new pods and cause service downtime.

This PR let each executor type adopts existing resources before starting the executor
API services, and so the alive connections won't experience failure. However, the requests
send to the function that doesn't have alive function pods will still fail due to the
executor is in bootstrapping.
This commit is contained in:
Ta-Ching Chen
2019-11-27 23:08:45 +08:00
committed by GitHub
parent 86446c8879
commit 1ad7ac2dcf
11 changed files with 437 additions and 115 deletions
+3 -13
View File
@@ -17,7 +17,6 @@ limitations under the License.
package executor
import (
"context"
"encoding/json"
"fmt"
"io/ioutil"
@@ -164,8 +163,8 @@ func (executor *Executor) tapServices(w http.ResponseWriter, r *http.Request) {
err = et.TapService(svcHost)
if err != nil {
errs = multierror.Append(errs,
errors.Wrapf(err, "'%v' failed to tap function '%v/%v' with service url '%v'",
req.FnMetadata.Namespace, req.FnMetadata.Name, req.ServiceUrl, req.FnExecutorType))
errors.Wrapf(err, "'%v' failed to tap function '%v' in '%v' with service url '%v'",
req.FnMetadata.Name, req.FnMetadata.Namespace, req.ServiceUrl, req.FnExecutorType))
}
}
@@ -192,16 +191,7 @@ func (executor *Executor) GetHandler() http.Handler {
}
func (executor *Executor) Serve(port int) {
executor.logger.Info("starting executor", zap.Int("port", port))
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
for _, et := range executor.executorTypes {
et.Run(ctx)
}
executor.cms.Run(ctx)
executor.logger.Info("starting executor API", zap.Int("port", port))
address := fmt.Sprintf(":%v", port)
err := http.ListenAndServe(address, &ochttp.Handler{
Handler: executor.GetHandler(),
+1 -1
View File
@@ -83,7 +83,7 @@ func initConfigmapController(logger *zap.Logger, fissionClient *crd.FissionClien
}
funcs, err := getConfigmapRelatedFuncs(logger, &newCm.ObjectMeta, fissionClient)
if err != nil {
logger.Error("Failed to get functions related to secret", zap.String("secret_name", newCm.ObjectMeta.Name), zap.String("secret_namespace", newCm.ObjectMeta.Namespace))
logger.Error("Failed to get functions related to configmap", zap.String("configmap_name", newCm.ObjectMeta.Name), zap.String("configmap_namespace", newCm.ObjectMeta.Namespace))
}
recyclePods(logger, funcs, types)
}
+18 -5
View File
@@ -217,32 +217,45 @@ func StartExecutor(logger *zap.Logger, functionNamespace string, envBuilderNames
}
restClient := fissionClient.GetCrdClient()
executorInstanceID := strings.ToLower(uniuri.NewLen(8))
poolID := strings.ToLower(uniuri.NewLen(8))
reaper.CleanupOldExecutorObjects(logger, kubernetesClient, poolID)
go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30)
logger.Info("Starting executor", zap.String("instanceID", executorInstanceID))
gpm := poolmgr.MakeGenericPoolManager(
logger,
fissionClient, kubernetesClient,
functionNamespace, fetcherConfig, poolID)
functionNamespace, fetcherConfig, executorInstanceID)
ndm := newdeploy.MakeNewDeploy(
logger,
fissionClient, kubernetesClient, restClient,
functionNamespace, fetcherConfig, poolID)
functionNamespace, fetcherConfig, executorInstanceID)
executorTypes := make(map[fv1.ExecutorType]executortype.ExecutorType)
executorTypes[gpm.GetTypeName()] = gpm
executorTypes[ndm.GetTypeName()] = ndm
wg := &sync.WaitGroup{}
for _, et := range executorTypes {
wg.Add(1)
go func(et executortype.ExecutorType) {
defer wg.Done()
et.AdoptOrphanResources()
et.Run(context.Background())
}(et)
}
wg.Wait()
cms := cms.MakeConfigSecretController(logger, fissionClient, kubernetesClient, executorTypes)
cms.Run(context.Background())
api, err := MakeExecutor(logger, cms, fissionClient, executorTypes)
if err != nil {
return err
}
go reaper.CleanupRoleBindings(logger, kubernetesClient, fissionClient, functionNamespace, envBuilderNamespace, time.Minute*30)
go reaper.CleanupOldExecutorObjects(logger, kubernetesClient, executorInstanceID)
go api.Serve(port)
go serveMetric(logger)
@@ -50,4 +50,7 @@ type ExecutorType interface {
// RefreshFuncPods refreshes function pods if the secrets/configmaps pods reference to get updated.
RefreshFuncPods(*zap.Logger, fv1.Function) error
// AdoptOrphanResources adopts existing resources created by the deleted executor.
AdoptOrphanResources()
}
@@ -29,6 +29,7 @@ import (
k8s_err "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/intstr"
fv1 "github.com/fission/fission/pkg/apis/fission.io/v1"
@@ -43,7 +44,7 @@ const (
)
func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Environment,
deployName string, deployLabels map[string]string, deployNamespace string, firstcreate bool) (*appsv1.Deployment, error) {
deployName string, deployLabels map[string]string, deployAnnotations map[string]string, deployNamespace string, firstcreate bool) (*appsv1.Deployment, error) {
minScale := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale)
specializationTimeout := int(fn.Spec.InvokeStrategy.ExecutionStrategy.SpecializationTimeout)
@@ -58,13 +59,23 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
existingDepl, err := deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Get(deployName, metav1.GetOptions{})
if err == nil {
if waitForDeploy {
err = deploy.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale)
if existingDepl.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID {
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID)
existingDepl, err = deploy.kubernetesClient.AppsV1().Deployments(deployNamespace).Patch(deployName, k8sTypes.StrategicMergePatchType, []byte(patch))
if err != nil {
deploy.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.Metadata.Name))
deploy.logger.Warn("error patching executor instance ID of deploy", zap.Error(err),
zap.String("deploy", deployName), zap.String("ns", deployNamespace))
return nil, err
}
}
if waitForDeploy {
if *existingDepl.Spec.Replicas < minScale {
err = deploy.scaleDeployment(existingDepl.Namespace, existingDepl.Name, minScale)
if err != nil {
deploy.logger.Error("error scaling up function deployment", zap.Error(err), zap.String("function", fn.Metadata.Name))
return nil, err
}
}
if existingDepl.Status.AvailableReplicas < minScale {
existingDepl, err = deploy.waitForDeploy(existingDepl, minScale, specializationTimeout)
}
@@ -76,7 +87,7 @@ func (deploy *NewDeploy) createOrGetDeployment(fn *fv1.Function, env *fv1.Enviro
return nil, err
}
deployment, err := deploy.getDeploymentSpec(fn, env, deployName, deployNamespace, deployLabels)
deployment, err := deploy.getDeploymentSpec(fn, env, deployName, deployNamespace, deployLabels, deployAnnotations)
if err != nil {
return nil, err
}
@@ -158,7 +169,7 @@ func (deploy *NewDeploy) deleteDeployment(ns string, name string) error {
}
func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environment,
deployName string, deployNamespace string, deployLabels map[string]string) (*appsv1.Deployment, error) {
deployName string, deployNamespace string, deployLabels map[string]string, deployAnnotations map[string]string) (*appsv1.Deployment, error) {
replicas := int32(fn.Spec.InvokeStrategy.ExecutionStrategy.MinScale)
@@ -174,6 +185,9 @@ func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environmen
if deploy.useIstio && env.Spec.AllowAccessToExternalNetwork {
podAnnotations["sidecar.istio.io/inject"] = "false"
}
for k, v := range deployAnnotations {
podAnnotations[k] = v
}
resources := deploy.getResources(env, fn)
// Set maxUnavailable and maxSurge to 20% is because we want
@@ -243,8 +257,9 @@ func (deploy *NewDeploy) getDeploymentSpec(fn *fv1.Function, env *fv1.Environmen
deployment := &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: deployName,
Labels: deployLabels,
Name: deployName,
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: appsv1.DeploymentSpec{
Replicas: &replicas,
@@ -319,7 +334,10 @@ func (deploy *NewDeploy) getResources(env *fv1.Environment, fn *fv1.Function) ap
return resources
}
func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy, depl *appsv1.Deployment) (*asv1.HorizontalPodAutoscaler, error) {
func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.ExecutionStrategy, depl *appsv1.Deployment, deployLabels map[string]string, deployAnnotations map[string]string) (*asv1.HorizontalPodAutoscaler, error) {
if depl == nil {
return nil, errors.New("failed to create HPA, found empty deployment")
}
minRepl := int32(execStrategy.MinScale)
if minRepl == 0 {
@@ -333,18 +351,22 @@ func (deploy *NewDeploy) createOrGetHpa(hpaName string, execStrategy *fv1.Execut
existingHpa, err := deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Get(hpaName, metav1.GetOptions{})
if err == nil {
if existingHpa.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID {
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID)
existingHpa, err = deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(depl.ObjectMeta.Namespace).Patch(hpaName, k8sTypes.StrategicMergePatchType, []byte(patch))
if err != nil {
deploy.logger.Warn("error patching executor instance ID of HPA", zap.Error(err),
zap.String("HPA", hpaName), zap.String("ns", depl.ObjectMeta.Namespace))
return nil, err
}
}
return existingHpa, err
}
if depl == nil {
return nil, errors.New("failed to create HPA, found empty deployment")
}
if err != nil && k8s_err.IsNotFound(err) {
} else if k8s_err.IsNotFound(err) {
hpa := asv1.HorizontalPodAutoscaler{
ObjectMeta: metav1.ObjectMeta{
Name: hpaName,
Labels: depl.Labels,
Name: hpaName,
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: asv1.HorizontalPodAutoscalerSpec{
ScaleTargetRef: asv1.CrossVersionObjectReference{
@@ -382,15 +404,25 @@ func (deploy *NewDeploy) deleteHpa(ns string, name string) error {
return deploy.kubernetesClient.AutoscalingV1().HorizontalPodAutoscalers(ns).Delete(name, &metav1.DeleteOptions{})
}
func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
func (deploy *NewDeploy) createOrGetSvc(deployLabels map[string]string, deployAnnotations map[string]string, svcName string, svcNamespace string) (*apiv1.Service, error) {
existingSvc, err := deploy.kubernetesClient.CoreV1().Services(svcNamespace).Get(svcName, metav1.GetOptions{})
if err == nil {
if existingSvc.Labels[fv1.EXECUTOR_INSTANCEID_LABEL] != deploy.instanceID {
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID)
existingSvc, err = deploy.kubernetesClient.CoreV1().Services(svcNamespace).Patch(svcName, k8sTypes.StrategicMergePatchType, []byte(patch))
if err != nil {
deploy.logger.Warn("error patching executor instance ID of service", zap.Error(err),
zap.String("service", svcName), zap.String("ns", svcNamespace))
return nil, err
}
}
return existingSvc, err
} else if k8s_err.IsNotFound(err) {
service := &apiv1.Service{
ObjectMeta: metav1.ObjectMeta{
Name: svcName,
Labels: deployLabels,
Name: svcName,
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: apiv1.ServiceSpec{
Ports: []apiv1.ServicePort{
@@ -19,9 +19,11 @@ package newdeploy
import (
"context"
"fmt"
"math/rand"
"os"
"strconv"
"strings"
"sync"
"time"
multierror "github.com/hashicorp/go-multierror"
@@ -171,7 +173,9 @@ func (deploy *NewDeploy) IsValid(fsvc *fscache.FuncSvc) bool {
_, err := deploy.kubernetesClient.CoreV1().Services(service[1]).Get(service[0], metav1.GetOptions{})
if err != nil {
deploy.logger.Error("error validating function service address", zap.String("function", fsvc.Function.Name), zap.Error(err))
if !k8sErrs.IsNotFound(err) {
deploy.logger.Error("error validating function service address", zap.String("function", fsvc.Function.Name), zap.Error(err))
}
return false
}
@@ -184,7 +188,9 @@ func (deploy *NewDeploy) IsValid(fsvc *fscache.FuncSvc) bool {
currentDeploy, err := deploy.kubernetesClient.AppsV1().
Deployments(deployObj.Namespace).Get(deployObj.Name, metav1.GetOptions{})
if err != nil {
deploy.logger.Error("error validating function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
if !k8sErrs.IsNotFound(err) {
deploy.logger.Error("error validating function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
}
return false
}
@@ -235,6 +241,78 @@ func (deploy *NewDeploy) RefreshFuncPods(logger *zap.Logger, f fv1.Function) err
return nil
}
func (deploy *NewDeploy) AdoptOrphanResources() {
l := map[string]string{
types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy),
}
podList, err := deploy.kubernetesClient.CoreV1().Pods(metav1.NamespaceAll).List(metav1.ListOptions{
LabelSelector: labels.Set(l).AsSelector().String(),
})
if err != nil {
deploy.logger.Error("error getting pod list", zap.Error(err))
return
}
podWG := &sync.WaitGroup{}
for i := range podList.Items {
pod := &podList.Items[i]
if !utils.IsReadyPod(pod) {
continue
}
podWG.Add(1)
go func() {
defer podWG.Done()
// avoid too many requests arrive Kubernetes API server at the same time.
time.Sleep(time.Duration(rand.Intn(30)) * time.Millisecond)
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, deploy.instanceID)
pod, err = deploy.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch))
if err != nil {
// just log the error since it won't affect the function serving
deploy.logger.Warn("error patching executor instance ID of pod", zap.Error(err),
zap.String("pod", pod.Name), zap.String("ns", pod.Namespace))
return
}
deploy.logger.Info("adopt newdeploy function pod",
zap.String("pod", pod.Name), zap.Any("labels", pod.Labels), zap.Any("annotations", pod.Annotations))
}()
}
fnList, err := deploy.fissionClient.Functions(metav1.NamespaceAll).List(metav1.ListOptions{})
if err != nil {
deploy.logger.Error("error getting function list", zap.Error(err))
return
}
deployWG := &sync.WaitGroup{}
for i := range fnList.Items {
fn := &fnList.Items[i]
if fn.Spec.InvokeStrategy.ExecutionStrategy.ExecutorType == fv1.ExecutorTypeNewdeploy {
deployWG.Add(1)
go func() {
defer deployWG.Done()
_, err = deploy.fnCreate(fn, true)
if err != nil {
deploy.logger.Warn("failed to adopt resources for function", zap.Error(err))
return
}
deploy.logger.Info("adopt resources for function", zap.String("function", fn.Metadata.Name))
}()
}
}
podWG.Wait()
deployWG.Wait()
}
func (deploy *NewDeploy) initFuncController() (k8sCache.Store, k8sCache.Controller) {
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(deploy.crdClient, "functions", metav1.NamespaceAll, fields.Everything())
@@ -378,6 +456,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function, firstcreate bool) (*fscache.
objName := deploy.getObjName(fn)
deployLabels := deploy.getDeployLabels(fn.Metadata, env.Metadata)
deployAnnotations := deploy.getDeployAnnotations(fn.Metadata)
// to support backward compatibility, if the function was created in default ns, we fall back to creating the
// deployment of the function in fission-function ns
@@ -391,7 +470,7 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function, firstcreate bool) (*fscache.
// Since newdeploy waits for pods of deployment to be ready,
// change the order of kubeObject creation (create service first,
// then deployment) to take advantage of waiting time.
svc, err := deploy.createOrGetSvc(deployLabels, objName, ns)
svc, err := deploy.createOrGetSvc(deployLabels, deployAnnotations, objName, ns)
if err != nil {
deploy.logger.Error("error creating service", zap.Error(err), zap.String("service", objName))
go deploy.cleanupNewdeploy(ns, objName)
@@ -399,14 +478,14 @@ func (deploy *NewDeploy) fnCreate(fn *fv1.Function, firstcreate bool) (*fscache.
}
svcAddress := fmt.Sprintf("%v.%v", svc.Name, svc.Namespace)
depl, err := deploy.createOrGetDeployment(fn, env, objName, deployLabels, ns, firstcreate)
depl, err := deploy.createOrGetDeployment(fn, env, objName, deployLabels, deployAnnotations, ns, firstcreate)
if err != nil {
deploy.logger.Error("error creating deployment", zap.Error(err), zap.String("deployment", objName))
go deploy.cleanupNewdeploy(ns, objName)
return nil, errors.Wrapf(err, "error creating deployment %v", objName)
}
hpa, err := deploy.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl)
hpa, err := deploy.createOrGetHpa(objName, &fn.Spec.InvokeStrategy.ExecutionStrategy, depl, deployLabels, deployAnnotations)
if err != nil {
deploy.logger.Error("error creating HPA", zap.Error(err), zap.String("hpa", objName))
go deploy.cleanupNewdeploy(ns, objName)
@@ -607,7 +686,7 @@ func (deploy *NewDeploy) updateFuncDeployment(fn *fv1.Function, env *fv1.Environ
ns = fn.Metadata.Namespace
}
newDeployment, err := deploy.getDeploymentSpec(fn, env, fnObjName, ns, deployLabels)
newDeployment, err := deploy.getDeploymentSpec(fn, env, fnObjName, ns, deployLabels, deploy.getDeployAnnotations(fn.Metadata))
if err != nil {
deploy.updateStatus(fn, err, "failed to get new deployment spec while updating function")
return err
@@ -665,15 +744,21 @@ func (deploy *NewDeploy) getObjName(fn *fv1.Function) string {
}
func (deploy *NewDeploy) getDeployLabels(fnMeta metav1.ObjectMeta, envMeta metav1.ObjectMeta) map[string]string {
return map[string]string{
types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy),
types.ENVIRONMENT_NAME: envMeta.Name,
types.ENVIRONMENT_NAMESPACE: envMeta.Namespace,
types.ENVIRONMENT_UID: string(envMeta.UID),
types.FUNCTION_NAME: fnMeta.Name,
types.FUNCTION_NAMESPACE: fnMeta.Namespace,
types.FUNCTION_UID: string(fnMeta.UID),
}
}
func (deploy *NewDeploy) getDeployAnnotations(fnMeta metav1.ObjectMeta) map[string]string {
return map[string]string{
types.EXECUTOR_INSTANCEID_LABEL: deploy.instanceID,
types.EXECUTOR_TYPE: string(fv1.ExecutorTypeNewdeploy),
types.ENVIRONMENT_NAME: envMeta.Name,
types.ENVIRONMENT_NAMESPACE: envMeta.Namespace,
types.ENVIRONMENT_UID: string(envMeta.UID),
types.FUNCTION_NAME: fnMeta.Name,
types.FUNCTION_NAMESPACE: fnMeta.Namespace,
types.FUNCTION_UID: string(fnMeta.UID),
types.FUNCTION_RESOURCE_VERSION: fnMeta.ResourceVersion,
}
}
@@ -739,7 +824,7 @@ func (deploy *NewDeploy) idleObjectReaper() {
currentDeploy, err := deploy.kubernetesClient.AppsV1().
Deployments(deployObj.Namespace).Get(deployObj.Name, metav1.GetOptions{})
if err != nil {
deploy.logger.Error("error validating function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
deploy.logger.Error("error getting function deployment", zap.Error(err), zap.String("function", fsvc.Function.Name))
continue
}
+73 -22
View File
@@ -33,8 +33,10 @@ import (
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
apiv1 "k8s.io/api/core/v1"
k8sErrs "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/intstr"
"k8s.io/client-go/kubernetes"
@@ -64,9 +66,9 @@ type (
kubernetesClient *kubernetes.Clientset
fissionClient *crd.FissionClient
instanceId string // poolmgr instance id
labelsForPool map[string]string
requestChannel chan *choosePodRequest
fetcherConfig *fetcherConfig.Config
stopCh context.CancelFunc
}
// serialize the choosing of pods so that choices don't conflict
@@ -97,6 +99,8 @@ func MakeGenericPool(
gpLogger.Info("creating pool", zap.Any("environment", env.Metadata))
ctx, stopCh := context.WithCancel(context.Background())
// TODO: in general we need to provide the user a way to configure pools. Initial
// replicas, autoscaling params, various timeouts, etc.
gp := &GenericPool{
@@ -116,6 +120,7 @@ func MakeGenericPool(
instanceId: instanceId,
useSvc: false, // defaults off -- svc takes a second or more to become routable, slowing cold start
useIstio: enableIstio, // defaults off -- istio integration requires pod relabeling and it takes a second or more to become routable, slowing cold start
stopCh: stopCh,
}
gp.runtimeImagePullPolicy = utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY"))
@@ -127,7 +132,7 @@ func MakeGenericPool(
}
// Labels for generic deployment/RS/pods.
gp.labelsForPool = gp.getDeployLabels()
//gp.labelsForPool = gp.getDeployLabels()
// create the pool
err = gp.createPool()
@@ -136,31 +141,41 @@ func MakeGenericPool(
}
gpLogger.Info("deployment created", zap.Any("environment", env.Metadata))
go gp.choosePodService()
go gp.choosePodService(ctx)
return gp, nil
}
func (gp *GenericPool) getDeployLabels() map[string]string {
return map[string]string{
types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr),
types.ENVIRONMENT_NAME: gp.env.Metadata.Name,
types.ENVIRONMENT_NAMESPACE: gp.env.Metadata.Namespace,
types.ENVIRONMENT_UID: string(gp.env.Metadata.UID),
"managed": "true", // this allows us to easily find pods managed by the deployment
}
}
func (gp *GenericPool) getDeployAnnotations() map[string]string {
return map[string]string{
fv1.EXECUTOR_INSTANCEID_LABEL: gp.instanceId,
types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr),
types.ENVIRONMENT_NAME: gp.env.Metadata.Name,
types.ENVIRONMENT_NAMESPACE: gp.env.Metadata.Namespace,
types.ENVIRONMENT_UID: string(gp.env.Metadata.UID),
"managed": "true", // this allows us to easily find pods managed by the deployment
}
}
// choosePodService serializes the choosing of pods
func (gp *GenericPool) choosePodService() {
for req := range gp.requestChannel {
pod, err := gp._choosePod(req.newLabels)
if err != nil {
req.responseChannel <- &choosePodResponse{error: err}
continue
func (gp *GenericPool) choosePodService(ctx context.Context) {
for {
select {
case req := <-gp.requestChannel:
pod, err := gp._choosePod(req.newLabels)
if err != nil {
req.responseChannel <- &choosePodResponse{error: err}
continue
}
req.responseChannel <- &choosePodResponse{pod: pod}
case <-ctx.Done():
return
}
req.responseChannel <- &choosePodResponse{pod: pod}
}
}
@@ -330,12 +345,15 @@ func (gp *GenericPool) specializePod(ctx context.Context, pod *apiv1.Pod, fn *fv
// getPoolName returns a unique name of an environment
func (gp *GenericPool) getPoolName() string {
return strings.ToLower(fmt.Sprintf("poolmgr-%v-%v-%v", gp.env.Metadata.Name, gp.env.Metadata.Namespace, uniuri.NewLen(8)))
return strings.ToLower(fmt.Sprintf("poolmgr-%v-%v-%v", gp.env.Metadata.Name, gp.env.Metadata.Namespace, gp.env.Metadata.ResourceVersion))
}
// A pool is a deployment of generic containers for an env. This
// creates the pool but doesn't wait for any pods to be ready.
func (gp *GenericPool) createPool() error {
deployLabels := gp.getDeployLabels()
deployAnnotations := gp.getDeployAnnotations()
// Use long terminationGracePeriodSeconds for connection draining in case that
// pod still runs user functions.
gracePeriodSeconds := int64(6 * 60)
@@ -350,6 +368,9 @@ func (gp *GenericPool) createPool() error {
if gp.useIstio && gp.env.Spec.AllowAccessToExternalNetwork {
podAnnotations["sidecar.istio.io/inject"] = "false"
}
for k, v := range deployAnnotations {
podAnnotations[k] = v
}
container, err := util.MergeContainer(&apiv1.Container{
Name: gp.env.Metadata.Name,
@@ -391,7 +412,7 @@ func (gp *GenericPool) createPool() error {
pod := apiv1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{
Labels: gp.labelsForPool,
Labels: deployLabels,
Annotations: podAnnotations,
},
Spec: apiv1.PodSpec{
@@ -408,13 +429,14 @@ func (gp *GenericPool) createPool() error {
deployment := &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: gp.getPoolName(),
Labels: gp.labelsForPool,
Name: gp.getPoolName(),
Labels: deployLabels,
Annotations: deployAnnotations,
},
Spec: appsv1.DeploymentSpec{
Replicas: &gp.replicas,
Selector: &metav1.LabelSelector{
MatchLabels: gp.labelsForPool,
MatchLabels: deployLabels,
},
Template: pod,
},
@@ -434,11 +456,25 @@ func (gp *GenericPool) createPool() error {
deployment.Spec.Template.Spec = *newPodSpec
}
_, err = gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Get(deployment.Name, metav1.GetOptions{})
if err == nil {
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, fv1.EXECUTOR_INSTANCEID_LABEL, gp.instanceId)
depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Patch(deployment.Name, k8sTypes.StrategicMergePatchType, []byte(patch))
if err == nil {
gp.deployment = depl
return nil
}
} else if !k8sErrs.IsNotFound(err) {
gp.logger.Error("error getting deployment in kubernetes", zap.Error(err), zap.String("deployment", deployment.Name))
return err
}
depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Create(deployment)
if err != nil {
gp.logger.Error("error creating deployment in kubernetes", zap.Error(err), zap.String("deployment", deployment.Name))
return err
}
gp.deployment = depl
return nil
}
@@ -497,7 +533,7 @@ func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1.
Ports: []apiv1.ServicePort{
{
Protocol: apiv1.ProtocolTCP,
Port: 80,
Port: 8888,
TargetPort: intstr.FromInt(8888),
},
},
@@ -584,7 +620,7 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
// the fission router isn't in the same namespace, so return a
// namespace-qualified hostname
svcHost = fmt.Sprintf("%v.%v", svcName, gp.namespace)
svcHost = fmt.Sprintf("%v.%v:8888", svcName, gp.namespace)
} else if gp.useIstio {
svc := utils.GetFunctionIstioServiceName(fn.Metadata.Name, fn.Metadata.Namespace)
svcHost = fmt.Sprintf("%v.%v:8888", svc, gp.namespace)
@@ -592,6 +628,18 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
svcHost = fmt.Sprintf("%v:8888", pod.Status.PodIP)
}
// patch svc-host and resource version to the pod annotations for new executor to adopt the pod
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v","%v":"%v"}}}`,
types.ANNOTATION_SVC_HOST, svcHost, types.FUNCTION_RESOURCE_VERSION, fn.Metadata.ResourceVersion)
p, err := gp.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch))
if err != nil {
// just log the error since it won't affect the function serving
gp.logger.Warn("error patching svc-host to pod", zap.Error(err),
zap.String("pod", pod.Name), zap.String("ns", pod.Namespace))
} else {
pod = p
}
gp.logger.Info("specialized pod",
zap.String("pod", pod.ObjectMeta.Name),
zap.String("podNamespace", pod.ObjectMeta.Namespace),
@@ -634,10 +682,13 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
// destroys the pool -- the deployment, replicaset and pods
func (gp *GenericPool) destroy() error {
gp.stopCh()
deletePropagation := metav1.DeletePropagationBackground
delOpt := metav1.DeleteOptions{
PropagationPolicy: &deletePropagation,
}
err := gp.kubernetesClient.AppsV1().
Deployments(gp.namespace).Delete(gp.deployment.ObjectMeta.Name, &delOpt)
if err != nil {
+152 -5
View File
@@ -18,13 +18,16 @@ package poolmgr
import (
"context"
"fmt"
"math/rand"
"os"
"strconv"
"strings"
"sync"
"time"
"github.com/fission/fission/pkg/utils"
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
k8serrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
@@ -40,6 +43,7 @@ import (
"github.com/fission/fission/pkg/executor/reaper"
fetcherConfig "github.com/fission/fission/pkg/fetcher/config"
"github.com/fission/fission/pkg/types"
"github.com/fission/fission/pkg/utils"
)
var _ executortype.ExecutorType = &GenericPoolManager{}
@@ -110,6 +114,7 @@ func MakeGenericPoolManager(
idlePodReapTime: 2 * time.Minute,
fetcherConfig: fetcherConfig,
}
go gpm.service()
go gpm.eagerPoolCreator()
@@ -238,6 +243,140 @@ func (gpm *GenericPoolManager) RefreshFuncPods(logger *zap.Logger, f fv1.Functio
return nil
}
func (gpm *GenericPoolManager) AdoptOrphanResources() {
envs, err := gpm.fissionClient.Environments(metav1.NamespaceAll).List(metav1.ListOptions{})
if err != nil {
gpm.logger.Error("error getting environment list", zap.Error(err))
return
}
envMap := make(map[string]fv1.Environment, len(envs.Items))
envWG := &sync.WaitGroup{}
for i := range envs.Items {
env := envs.Items[i]
if gpm.getEnvPoolsize(&env) > 0 {
envWG.Add(1)
go func() {
defer envWG.Done()
_, err := gpm.getPool(&env)
if err != nil {
gpm.logger.Error("adopt pool failed", zap.Error(err))
}
}()
}
// create environment map for later use
key := fmt.Sprintf("%v/%v", env.Metadata.Namespace, env.Metadata.Name)
envMap[key] = env
}
l := map[string]string{
types.EXECUTOR_TYPE: string(fv1.ExecutorTypePoolmgr),
}
podList, err := gpm.kubernetesClient.CoreV1().Pods(metav1.NamespaceAll).List(metav1.ListOptions{
LabelSelector: labels.Set(l).AsSelector().String(),
})
if err != nil {
gpm.logger.Error("error getting pod list", zap.Error(err))
return
}
podWG := &sync.WaitGroup{}
for i := range podList.Items {
pod := &podList.Items[i]
if !utils.IsReadyPod(pod) {
continue
}
podWG.Add(1)
go func() {
defer podWG.Done()
// avoid too many requests arrive Kubernetes API server at the same time.
time.Sleep(time.Duration(rand.Intn(30)) * time.Millisecond)
patch := fmt.Sprintf(`{"metadata":{"annotations":{"%v":"%v"}}}`, types.EXECUTOR_INSTANCEID_LABEL, gpm.instanceId)
pod, err = gpm.kubernetesClient.CoreV1().Pods(pod.Namespace).Patch(pod.Name, k8sTypes.StrategicMergePatchType, []byte(patch))
if err != nil {
// just log the error since it won't affect the function serving
gpm.logger.Warn("error patching executor instance ID of pod", zap.Error(err),
zap.String("pod", pod.Name), zap.String("ns", pod.Namespace))
return
}
// for unspecialized pod, we only update its annotations
if pod.Labels["managed"] == "true" {
return
}
fnName, ok1 := pod.Labels[types.FUNCTION_NAME]
fnNS, ok2 := pod.Labels[types.FUNCTION_NAMESPACE]
fnUID, ok3 := pod.Labels[types.FUNCTION_UID]
fnRV, ok4 := pod.Annotations[types.FUNCTION_RESOURCE_VERSION]
envName, ok5 := pod.Labels[types.ENVIRONMENT_NAME]
envNS, ok6 := pod.Labels[types.ENVIRONMENT_NAMESPACE]
svcHost, ok7 := pod.Annotations[types.ANNOTATION_SVC_HOST]
env, ok8 := envMap[fmt.Sprintf("%v/%v", envNS, envName)]
if !(ok1 && ok2 && ok3 && ok4 && ok5 && ok6 && ok7 && ok8) {
gpm.logger.Warn("failed to adopt pod for function due to lack necessary information",
zap.String("pod", pod.Name), zap.Any("labels", pod.Labels), zap.Any("annotations", pod.Annotations),
zap.String("env", env.Metadata.Name))
return
}
fsvc := fscache.FuncSvc{
Name: pod.Name,
Function: &metav1.ObjectMeta{
Name: fnName,
Namespace: fnNS,
UID: k8sTypes.UID(fnUID),
ResourceVersion: fnRV,
},
Environment: &env,
Address: svcHost,
KubernetesObjects: []apiv1.ObjectReference{
{
Kind: "pod",
Name: pod.Name,
APIVersion: pod.TypeMeta.APIVersion,
Namespace: pod.ObjectMeta.Namespace,
ResourceVersion: pod.ObjectMeta.ResourceVersion,
UID: pod.ObjectMeta.UID,
},
},
Executor: fv1.ExecutorTypePoolmgr,
Ctime: time.Now(),
Atime: time.Now(),
}
// If fsvc already exists we just skip the duplicate one. And let reaper to recycle the duplicate pod.
// This is for the case that there are multiple function pods for the same function due to unknown reason.
_, err := gpm.fsCache.GetByFunction(fsvc.Function)
if err == nil {
return
}
_, err = gpm.fsCache.Add(fsvc)
if err != nil {
gpm.logger.Warn("failed to adopt pod for function", zap.Error(err), zap.String("pod", pod.Name))
return
}
gpm.logger.Info("adopt function pod",
zap.String("pod", pod.Name), zap.Any("labels", pod.Labels), zap.Any("annotations", pod.Annotations))
}()
}
envWG.Wait()
podWG.Wait()
}
func (gpm *GenericPoolManager) service() {
for {
req := <-gpm.requestChannel
@@ -353,19 +492,27 @@ func (gpm *GenericPoolManager) eagerPoolCreator() {
// creating pools for envs that are actually used by functions. Also we might want
// to keep these eagerly created pools smaller than the ones created when there are
// actual function calls.
wg := &sync.WaitGroup{}
for i := range envs.Items {
env := envs.Items[i]
// Create pool only if poolsize greater than zero
if gpm.getEnvPoolsize(&env) > 0 {
_, err := gpm.getPool(&envs.Items[i])
if err != nil {
gpm.logger.Error("eager-create pool failed", zap.Error(err))
}
wg.Add(1)
go func() {
defer wg.Done()
_, err := gpm.getPool(&env)
if err != nil {
gpm.logger.Error("eager-create pool failed", zap.Error(err))
}
}()
}
}
// Clean up pools whose env was deleted
gpm.cleanupPools(envs.Items)
wg.Wait()
time.Sleep(pollSleep)
}
}
+20 -24
View File
@@ -38,16 +38,20 @@ var (
// CleanupOldExecutorObjects cleans up resources created by old executor instances
func CleanupOldExecutorObjects(logger *zap.Logger, kubernetesClient *kubernetes.Clientset, instanceId string) {
go func() {
err := cleanup(logger, kubernetesClient, instanceId)
if err != nil {
// TODO retry reaper; logged and ignored for now
logger.Error("Failed to cleanup old executor objects", zap.Error(err))
}
}()
err := cleanup(logger, kubernetesClient, instanceId)
if err != nil {
// TODO retry reaper; logged and ignored for now
logger.Error("Failed to cleanup old executor objects", zap.Error(err))
}
}
func cleanup(logger *zap.Logger, client *kubernetes.Clientset, instanceId string) error {
// Pods might still 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, client, instanceId)
if err != nil {
@@ -67,13 +71,6 @@ func cleanup(logger *zap.Logger, client *kubernetes.Clientset, instanceId string
return err
}
// Pods might still 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 = cleanupPods(logger, client, instanceId)
if err != nil {
return err
@@ -121,7 +118,7 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan
return err
}
for _, dep := range deploymentList.Items {
id, ok := dep.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
id, ok := dep.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL]
if ok && id != instanceId {
logger.Debug("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name))
err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt)
@@ -134,7 +131,7 @@ func cleanupDeployments(logger *zap.Logger, client *kubernetes.Clientset, instan
// ignore err
}
// Backward compatibility with older label name
pid, pok := dep.ObjectMeta.Labels[types.POOLMGR_INSTANCEID_LABEL]
pid, pok := dep.ObjectMeta.Annotations[types.POOLMGR_INSTANCEID_LABEL]
if pok && pid != instanceId {
logger.Debug("cleaning up deployment", zap.String("deployment", dep.ObjectMeta.Name))
err := client.AppsV1().Deployments(dep.ObjectMeta.Namespace).Delete(dep.ObjectMeta.Name, &delOpt)
@@ -156,7 +153,7 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st
return err
}
for _, pod := range podList.Items {
id, ok := pod.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
id, ok := pod.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL]
if ok && id != instanceId {
logger.Debug("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name))
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil)
@@ -169,7 +166,7 @@ func cleanupPods(logger *zap.Logger, client *kubernetes.Clientset, instanceId st
// ignore err
}
// Backward compatibility with older label name
pid, pok := pod.ObjectMeta.Labels[types.POOLMGR_INSTANCEID_LABEL]
pid, pok := pod.ObjectMeta.Annotations[types.POOLMGR_INSTANCEID_LABEL]
if pok && pid != instanceId {
logger.Debug("cleaning up pod", zap.String("pod", pod.ObjectMeta.Name))
err := client.CoreV1().Pods(pod.ObjectMeta.Namespace).Delete(pod.ObjectMeta.Name, nil)
@@ -190,7 +187,7 @@ func cleanupServices(logger *zap.Logger, client *kubernetes.Clientset, instanceI
return err
}
for _, svc := range svcList.Items {
id, ok := svc.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
id, ok := svc.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL]
if ok && id != instanceId {
logger.Debug("cleaning up service", zap.String("service", svc.ObjectMeta.Name))
err := client.CoreV1().Services(svc.ObjectMeta.Namespace).Delete(svc.ObjectMeta.Name, nil)
@@ -213,7 +210,7 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str
}
for _, hpa := range hpaList.Items {
id, ok := hpa.ObjectMeta.Labels[types.EXECUTOR_INSTANCEID_LABEL]
id, ok := hpa.ObjectMeta.Annotations[types.EXECUTOR_INSTANCEID_LABEL]
if ok && id != instanceId {
logger.Debug("cleaning up HPA", zap.String("hpa", hpa.ObjectMeta.Name))
err := client.AutoscalingV1().HorizontalPodAutoscalers(hpa.ObjectMeta.Namespace).Delete(hpa.ObjectMeta.Name, nil)
@@ -225,7 +222,6 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str
}
// ignore err
}
}
return nil
@@ -235,6 +231,9 @@ func cleanupHpa(logger *zap.Logger, client *kubernetes.Clientset, instanceId str
// deletes the rolebindings completely if there are no Service Accounts in a rolebinding object.
func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissionClient *crd.FissionClient, functionNs, envBuilderNs string, cleanupRoleBindingInterval time.Duration) {
for {
// some sleep before the next reaper iteration
time.Sleep(cleanupRoleBindingInterval)
logger.Debug("starting cleanupRoleBindings cycle")
// get all rolebindings ( just to be efficient, one call to kubernetes )
rbList, err := client.RbacV1beta1().RoleBindings(meta_v1.NamespaceAll).List(meta_v1.ListOptions{})
@@ -363,8 +362,5 @@ func CleanupRoleBindings(logger *zap.Logger, client *kubernetes.Clientset, fissi
}
}
}
// some sleep before the next reaper iteration
time.Sleep(cleanupRoleBindingInterval)
}
}
+12 -7
View File
@@ -121,13 +121,18 @@ const (
// executor kubernetes object label key
const (
ENVIRONMENT_NAMESPACE = "environmentNamespace"
ENVIRONMENT_NAME = "environmentName"
ENVIRONMENT_UID = "environmentUid"
FUNCTION_NAMESPACE = "functionNamespace"
FUNCTION_NAME = "functionName"
FUNCTION_UID = "functionUid"
EXECUTOR_TYPE = "executorType"
ENVIRONMENT_NAMESPACE = "environmentNamespace"
ENVIRONMENT_NAME = "environmentName"
ENVIRONMENT_UID = "environmentUid"
FUNCTION_NAMESPACE = "functionNamespace"
FUNCTION_NAME = "functionName"
FUNCTION_UID = "functionUid"
FUNCTION_RESOURCE_VERSION = "functionResourceVersion"
EXECUTOR_TYPE = "executorType"
)
const (
ANNOTATION_SVC_HOST = "svcHost"
)
const (