Readypod optimization in executor (#1860)

A ready pod which can be specialized was fetched for every function earlier, this has been changed to a queue and cache implementation in client-go to improve performance.
This commit is contained in:
Rahul Bhati
2020-11-19 15:23:02 +05:30
committed by GitHub
parent aaf9f18d93
commit bb9f4136f5
8 changed files with 128 additions and 180 deletions
+1 -1
View File
@@ -1,4 +1,4 @@
ARG GO_VERSION=1.9.2
ARG GO_VERSION=1.13
FROM tensorflow/serving as serving
RUN apt update && apt install -y ca-certificates && rm -rf /var/lib/apt/lists/*
+1 -2
View File
@@ -53,6 +53,7 @@ require (
github.com/opencontainers/image-spec v1.0.1 // indirect
github.com/opencontainers/runc v0.1.1 // indirect
github.com/ory/dockertest v3.3.5+incompatible
github.com/pierrec/lz4 v2.0.5+incompatible // indirect
github.com/pkg/errors v0.9.1
github.com/prometheus/client_golang v1.0.0
github.com/prometheus/common v0.4.1
@@ -65,8 +66,6 @@ require (
github.com/ulikunitz/xz v0.0.0-20180703112113-636d36a76670 // indirect
github.com/wcharczuk/go-chart v2.0.1+incompatible
go.opencensus.io v0.22.0
go.uber.org/atomic v1.3.2 // indirect
go.uber.org/multierr v1.1.0 // indirect
go.uber.org/zap v1.9.1
golang.org/x/image v0.0.0-20190618124811-92942e4437e2 // indirect
golang.org/x/net v0.0.0-20200202094626-16171245cfb2
+2 -4
View File
@@ -445,12 +445,10 @@ go.opencensus.io v0.20.1/go.mod h1:6WKK9ahsWS3RSO+PY9ZHZUfv2irvY6gN279GOPZjmmk=
go.opencensus.io v0.21.0/go.mod h1:mSImk1erAIZhrmZN+AvHh14ztQfjbGwt4TtuofqLduU=
go.opencensus.io v0.22.0 h1:C9hSCOW830chIVkdja34wa6Ky+IzWllkUinR+BtRZd4=
go.opencensus.io v0.22.0/go.mod h1:+kGneAE2xo2IficOXnaByMWTGM9T73dGwxeWcUqIpI8=
go.uber.org/atomic v0.0.0-20181018215023-8dc6146f7569 h1:nSQar3Y0E3VQF/VdZ8PTAilaXpER+d7ypdABCrpwMdg=
go.uber.org/atomic v0.0.0-20181018215023-8dc6146f7569/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
go.uber.org/atomic v1.3.2 h1:2Oa65PReHzfn29GpvgsYwloV9AVFHPDk8tYxt2c2tr4=
go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
go.uber.org/multierr v0.0.0-20180122172545-ddea229ff1df h1:shvkWr0NAZkg4nPuE3XrKP0VuBPijjk3TfX6Y6acFNg=
go.uber.org/multierr v0.0.0-20180122172545-ddea229ff1df/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0=
go.uber.org/multierr v1.1.0 h1:HoEmRHQPVSqub6w2z2d2EOVs2fjyFRGyofhKuyDq0QI=
go.uber.org/multierr v1.1.0/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0=
go.uber.org/zap v0.0.0-20180814183419-67bc79d13d15/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q=
go.uber.org/zap v1.9.1 h1:XCJQEf3W6eZaVwhRBof6ImoYGJSITeKWsyeh3HFu/5o=
go.uber.org/zap v1.9.1/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q=
+1 -1
View File
@@ -77,7 +77,7 @@ func (executor *Executor) getServiceForFunctionApi(w http.ResponseWriter, r *htt
executor.logger.Debug("setting concurrency to 5")
}
if t == fv1.ExecutorTypePoolmgr && et.GetTotalAvailable(fn) >= conncurrency {
errMsg := fmt.Sprintf("max concurrency reached for %v. All %v instance are active", fn.ObjectMeta.Name, fn.Spec.Concurrency)
errMsg := fmt.Sprintf("max concurrency reached for %v. All %v instance are active", fn.ObjectMeta.Name, conncurrency)
executor.logger.Error("error occurred", zap.String("error", errMsg))
http.Error(w, errMsg, http.StatusTooManyRequests)
return
+78 -162
View File
@@ -27,7 +27,6 @@ import (
"github.com/dchest/uniuri"
"github.com/fission/fission/pkg/utils"
multierror "github.com/hashicorp/go-multierror"
"github.com/pkg/errors"
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
@@ -38,6 +37,8 @@ import (
k8sTypes "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/intstr"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/util/workqueue"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
@@ -49,34 +50,26 @@ import (
type (
GenericPool struct {
logger *zap.Logger
env *fv1.Environment
replicas int32 // num idle pods
deployment *appsv1.Deployment // kubernetes deployment
namespace string // namespace to keep our resources
functionNamespace string // fallback namespace for fission functions
podReadyTimeout time.Duration // timeout for generic pods to become ready
fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and podname
useSvc bool // create k8s service for specialized pods
useIstio bool
poolInstanceId string // small random string to uniquify pod names
runtimeImagePullPolicy apiv1.PullPolicy // pull policy for generic pool to created env deployment
kubernetesClient *kubernetes.Clientset
fissionClient *crd.FissionClient
instanceId string // poolmgr instance id
requestChannel chan *choosePodRequest
fetcherConfig *fetcherConfig.Config
stopCh context.CancelFunc
}
// serialize the choosing of pods so that choices don't conflict
choosePodRequest struct {
newLabels map[string]string
responseChannel chan *choosePodResponse
}
choosePodResponse struct {
pod *apiv1.Pod
error
logger *zap.Logger
env *fv1.Environment
replicas int32 // num idle pods
deployment *appsv1.Deployment // kubernetes deployment
namespace string // namespace to keep our resources
functionNamespace string // fallback namespace for fission functions
podReadyTimeout time.Duration // timeout for generic pods to become ready
fsCache *fscache.FunctionServiceCache // cache funcSvc's by function, address and podname
useSvc bool // create k8s service for specialized pods
useIstio bool
poolInstanceId string // small random string to uniquify pod names
runtimeImagePullPolicy apiv1.PullPolicy // pull policy for generic pool to created env deployment
kubernetesClient *kubernetes.Clientset
fissionClient *crd.FissionClient
instanceId string // poolmgr instance id
fetcherConfig *fetcherConfig.Config
stopReadyPodControllerCh chan struct{}
readyPodController cache.Controller
readyPodIndexer cache.Indexer
readyPodQueue workqueue.RateLimitingInterface
}
)
@@ -107,27 +100,24 @@ func MakeGenericPool(
gpLogger.Info("creating pool", zap.Any("environment", env.ObjectMeta))
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{
logger: gpLogger,
env: env,
replicas: initialReplicas, // TODO make this an env param instead?
requestChannel: make(chan *choosePodRequest),
fissionClient: fissionClient,
kubernetesClient: kubernetesClient,
namespace: namespace,
functionNamespace: functionNamespace,
podReadyTimeout: podReadyTimeout,
fsCache: fsCache,
poolInstanceId: uniuri.NewLen(8),
fetcherConfig: fetcherConfig,
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,
logger: gpLogger,
env: env,
replicas: initialReplicas, // TODO make this an env param instead?
fissionClient: fissionClient,
kubernetesClient: kubernetesClient,
namespace: namespace,
functionNamespace: functionNamespace,
podReadyTimeout: podReadyTimeout,
fsCache: fsCache,
poolInstanceId: uniuri.NewLen(8),
fetcherConfig: fetcherConfig,
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
stopReadyPodControllerCh: make(chan struct{}),
}
gp.runtimeImagePullPolicy = utils.GetImagePullPolicy(os.Getenv("RUNTIME_IMAGE_PULL_POLICY"))
@@ -148,7 +138,7 @@ func MakeGenericPool(
}
gpLogger.Info("deployment created", zap.Any("environment", env.ObjectMeta))
go gp.choosePodService(ctx)
go gp.startReadyPodController()
return gp, nil
}
@@ -169,87 +159,50 @@ func (gp *GenericPool) getDeployAnnotations() map[string]string {
}
}
// choosePodService serializes the choosing of pods
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
}
}
}
// choosePod picks a ready pod from the pool and relabels it, waiting if necessary.
// returns the pod API object.
func (gp *GenericPool) choosePod(newLabels map[string]string) (*apiv1.Pod, error) {
req := &choosePodRequest{
newLabels: newLabels,
responseChannel: make(chan *choosePodResponse),
}
gp.requestChannel <- req
resp := <-req.responseChannel
return resp.pod, resp.error
}
// _choosePod is called serially by choosePodService
func (gp *GenericPool) _choosePod(newLabels map[string]string) (*apiv1.Pod, error) {
// returns the key and pod API object.
func (gp *GenericPool) choosePod(newLabels map[string]string) (string, *apiv1.Pod, error) {
startTime := time.Now()
for {
// Retries took too long, error out.
if time.Since(startTime) > gp.podReadyTimeout {
gp.logger.Error("timed out waiting for pod", zap.Any("labels", newLabels), zap.Duration("timeout", gp.podReadyTimeout))
return nil, errors.New("timeout: waited too long to get a ready pod")
return "", nil, errors.New("timeout: waited too long to get a ready pod")
}
// Get pods; filter the ones that are ready
podList, err := gp.kubernetesClient.CoreV1().Pods(gp.namespace).List(
metav1.ListOptions{
FieldSelector: "status.phase=Running",
LabelSelector: labels.Set(
gp.deployment.Spec.Selector.MatchLabels).AsSelector().String(),
})
if err != nil {
return nil, err
}
readyPods := make([]*apiv1.Pod, 0, len(podList.Items))
for i := range podList.Items {
pod := podList.Items[i]
var chosenPod *apiv1.Pod
var key string
// Ignore not ready pod here
if !utils.IsReadyPod(&pod) {
if gp.readyPodQueue.Len() > 0 {
item, quit := gp.readyPodQueue.Get()
if quit {
gp.logger.Error("readypod controller is not running")
return "", nil, errors.New("readypod controller is not running")
}
key := item.(string)
obj, exists, err := gp.readyPodIndexer.GetByKey(key)
if err != nil {
gp.logger.Error("fetching object from store failed", zap.String("key", key), zap.Error(err))
return "", nil, err
}
if !exists {
gp.logger.Warn("pod deleted from store", zap.String("pod", key))
continue
}
// add it to the list of ready pods
readyPods = append(readyPods, &pod)
break
}
gp.logger.Info("found ready pods",
zap.Any("labels", newLabels),
zap.Int("ready_count", len(readyPods)),
zap.Int("total", len(podList.Items)))
// If there are no ready pods, wait and retry.
if len(readyPods) == 0 {
err = gp.waitForReadyPod()
if err != nil {
return nil, err
if !utils.IsReadyPod(obj.(*apiv1.Pod)) {
continue
}
chosenPod = obj.(*apiv1.Pod).DeepCopy()
} else {
// Wait for pods to get ready and retry
gp.logger.Info("waiting for ready pods")
time.Sleep(1000 * time.Millisecond)
continue
}
// Pick a ready pod. For now just choose randomly;
// ideally we'd care about which node it's running on,
// and make a good scheduling decision.
chosenPod := readyPods[0]
if gp.env.Spec.AllowedFunctionsPerContainer != fv1.AllowedFunctionsPerContainerInfinite {
// Relabel. If the pod already got picked and
// modified, this should fail; in that case just
@@ -274,13 +227,13 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*apiv1.Pod, erro
// So we have to check both of them to ensure the patch success.
for k, v := range newLabels {
if newPod.Labels[k] != v {
return nil, errors.Errorf("value of necessary labels '%v' mismatch: want '%v', get '%v'",
return "", nil, errors.Errorf("value of necessary labels '%v' mismatch: want '%v', get '%v'",
k, v, newPod.Labels[k])
}
}
for k, v := range annotations {
if newPod.Annotations[k] != v {
return nil, errors.Errorf("value of necessary annotations '%v' mismatch: want '%v', get '%v'",
return "", nil, errors.Errorf("value of necessary annotations '%v' mismatch: want '%v', get '%v'",
k, v, newPod.Annotations[k])
}
}
@@ -289,7 +242,7 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*apiv1.Pod, erro
gp.logger.Info("chose pod", zap.Any("labels", newLabels),
zap.String("pod", chosenPod.Name), zap.Duration("elapsed_time", time.Since(startTime)))
return chosenPod, nil
return key, chosenPod, nil
}
}
@@ -522,50 +475,6 @@ func (gp *GenericPool) createPool() error {
return nil
}
func (gp *GenericPool) waitForReadyPod() error {
startTime := time.Now()
for {
// TODO: for now we just poll; use a watch instead
depl, err := gp.kubernetesClient.AppsV1().Deployments(gp.namespace).Get(
gp.deployment.ObjectMeta.Name, metav1.GetOptions{})
if err != nil {
e := "error waiting for ready pod for deployment"
gp.logger.Error(e, zap.String("deployment", gp.deployment.ObjectMeta.Name), zap.String("namespace", gp.namespace))
return fmt.Errorf("%s %q in namespace %q", e, gp.deployment.ObjectMeta.Name, gp.namespace)
}
gp.deployment = depl
if gp.deployment.Status.AvailableReplicas > 0 {
return nil
}
if time.Since(startTime) > gp.podReadyTimeout {
podList, err := gp.kubernetesClient.CoreV1().Pods(gp.namespace).List(metav1.ListOptions{
LabelSelector: labels.Set(
gp.deployment.Spec.Selector.MatchLabels).AsSelector().String(),
})
if err != nil {
gp.logger.Error("error getting pod list after timeout waiting for ready pod", zap.Error(err))
}
// Since even single pod is not ready, choosing the first pod to inspect is a good approximation. In future this can be done better
pod := podList.Items[0]
errs := &multierror.Error{}
for _, cStatus := range pod.Status.ContainerStatuses {
if !cStatus.Ready {
errs = multierror.Append(errs, errors.New(fmt.Sprintf("%v: %v", cStatus.State.Waiting.Reason, cStatus.State.Waiting.Message)))
}
}
if errs.ErrorOrNil() != nil {
return errors.Wrapf(errs, "Timeout: waited too long for pod of deployment %v in namespace %v to be ready",
gp.deployment.ObjectMeta.Name, gp.namespace)
}
return nil
}
time.Sleep(1000 * time.Millisecond)
}
}
func (gp *GenericPool) createSvc(name string, labels map[string]string) (*apiv1.Service, error) {
service := apiv1.Service{
ObjectMeta: metav1.ObjectMeta{
@@ -633,13 +542,20 @@ func (gp *GenericPool) getFuncSvc(ctx context.Context, fn *fv1.Function) (*fscac
}
}
pod, err := gp.choosePod(funcLabels)
key, pod, err := gp.choosePod(funcLabels)
if err != nil {
return nil, err
}
defer gp.readyPodQueue.Done(key)
err = gp.specializePod(ctx, pod, fn)
if err != nil {
if gp.readyPodQueue.NumRequeues(key) < 5 {
gp.readyPodQueue.AddRateLimited(key)
} else {
gp.readyPodQueue.Forget(key)
}
gp.scheduleDeletePod(pod.ObjectMeta.Name)
return nil, err
}
@@ -723,7 +639,7 @@ 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()
close(gp.stopReadyPodControllerCh)
deletePropagation := metav1.DeletePropagationBackground
delOpt := metav1.DeleteOptions{
@@ -0,0 +1,45 @@
package poolmgr
import (
"go.uber.org/zap"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/util/workqueue"
)
func (gp *GenericPool) startReadyPodController() {
// create the pod watcher to filter by labels
// Filtering pod by phase=Running. In some cases the pod can be in
// different sate than Running, for example Kubernetes sets a
// pod to Termination while k8s waits for the grace period of
// the pod, even if all the containers are in Ready state.
optionsModifier := func(options *metav1.ListOptions) {
options.LabelSelector = labels.Set(
gp.deployment.Spec.Selector.MatchLabels).AsSelector().String()
options.FieldSelector = "status.phase=Running"
}
readyPodWatcher := cache.NewFilteredListWatchFromClient(gp.kubernetesClient.CoreV1().RESTClient(), "pods", gp.namespace, optionsModifier)
gp.readyPodQueue = workqueue.NewRateLimitingQueue(workqueue.DefaultControllerRateLimiter())
gp.readyPodIndexer, gp.readyPodController = cache.NewIndexerInformer(readyPodWatcher, &apiv1.Pod{}, 0, cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
key, err := cache.MetaNamespaceKeyFunc(obj)
if err == nil {
gp.readyPodQueue.Add(key)
gp.logger.Debug("add func called", zap.String("key", key))
}
},
DeleteFunc: func(obj interface{}) {
key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
if err == nil {
gp.readyPodQueue.Forget(key)
gp.readyPodQueue.Done(key)
gp.logger.Debug("delete func called", zap.String("key", key))
}
},
}, cache.Indexers{})
go gp.readyPodController.Run(gp.stopReadyPodControllerCh)
gp.logger.Info("readyPod controller started", zap.String("env", gp.env.ObjectMeta.Name), zap.String("envID", string(gp.env.ObjectMeta.UID)))
}
-1
View File
@@ -148,7 +148,6 @@ func (w *fakeCloseReadCloser) RealClose() error {
func (roundTripper *RetryingRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
// set the timeout for transport context
roundTripper.addForwardedHostHeader(req)
transport := roundTripper.getDefaultTransport()
ocRoundTripper := &ochttp.Transport{Base: transport}
-9
View File
@@ -66,15 +66,6 @@ func IsReadyPod(pod *apiv1.Pod) bool {
return false
}
// pod is not in Running Phase. It can be in Pending,
// Succeeded, Failed, Unknown. In some cases the pod can be in
// different sate than Running, for example Kubernetes sets a
// pod to Termination while k8s waits for the grace period of
// the pod, even if all the containers are in Ready state.
if pod.Status.Phase != apiv1.PodRunning {
return false
}
// pod is in "Terminating" status if deletionTimestamp is not nil
// https://github.com/kubernetes/kubernetes/issues/61376
if pod.ObjectMeta.DeletionTimestamp != nil {