From dc1f17944ea8c51dbf9ec29a43acb93dafd0e59e Mon Sep 17 00:00:00 2001 From: Dharma Bellamkonda Date: Tue, 17 Aug 2021 19:49:41 -0600 Subject: [PATCH] Page list requests by default (#3823) Signed-off-by: Dharma Bellamkonda --- changelogs/unreleased/3823-dharmab | 1 + pkg/backup/backup.go | 6 +- pkg/backup/item_collector.go | 73 +++++++++++++++++++--- pkg/cmd/server/server.go | 13 +++- site/content/docs/main/backup-reference.md | 8 +++ 5 files changed, 88 insertions(+), 13 deletions(-) create mode 100644 changelogs/unreleased/3823-dharmab diff --git a/changelogs/unreleased/3823-dharmab b/changelogs/unreleased/3823-dharmab new file mode 100644 index 000000000..d441a228c --- /dev/null +++ b/changelogs/unreleased/3823-dharmab @@ -0,0 +1 @@ +Add --client-page-size flag to server to allow chunking Kubernetes API LIST calls across multiple requests on large clusters diff --git a/pkg/backup/backup.go b/pkg/backup/backup.go index e5a6b1438..746695e00 100644 --- a/pkg/backup/backup.go +++ b/pkg/backup/backup.go @@ -1,5 +1,5 @@ /* -Copyright 2020 the Velero contributors. +Copyright the Velero contributors. Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. @@ -73,6 +73,7 @@ type kubernetesBackupper struct { resticBackupperFactory restic.BackupperFactory resticTimeout time.Duration defaultVolumesToRestic bool + clientPageSize int } type resolvedAction struct { @@ -106,6 +107,7 @@ func NewKubernetesBackupper( resticBackupperFactory restic.BackupperFactory, resticTimeout time.Duration, defaultVolumesToRestic bool, + clientPageSize int, ) (Backupper, error) { return &kubernetesBackupper{ backupClient: backupClient, @@ -115,6 +117,7 @@ func NewKubernetesBackupper( resticBackupperFactory: resticBackupperFactory, resticTimeout: resticTimeout, defaultVolumesToRestic: defaultVolumesToRestic, + clientPageSize: clientPageSize, }, nil } @@ -272,6 +275,7 @@ func (kb *kubernetesBackupper) Backup(log logrus.FieldLogger, backupRequest *Req dynamicFactory: kb.dynamicFactory, cohabitatingResources: cohabitatingResources(), dir: tempDir, + pageSize: kb.clientPageSize, } items := collector.getAllItems() diff --git a/pkg/backup/item_collector.go b/pkg/backup/item_collector.go index acdf548ea..6029d8f97 100644 --- a/pkg/backup/item_collector.go +++ b/pkg/backup/item_collector.go @@ -1,5 +1,5 @@ /* -Copyright 2020 the Velero contributors. +Copyright the Velero contributors. Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. @@ -17,6 +17,7 @@ limitations under the License. package backup import ( + "context" "encoding/json" "fmt" "io/ioutil" @@ -25,10 +26,13 @@ import ( "github.com/pkg/errors" "github.com/sirupsen/logrus" + apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/tools/pager" "github.com/vmware-tanzu/velero/pkg/client" "github.com/vmware-tanzu/velero/pkg/discovery" @@ -45,6 +49,7 @@ type itemCollector struct { dynamicFactory client.DynamicFactory cohabitatingResources map[string]*cohabitatingResource dir string + pageSize int } type kubernetesResource struct { @@ -275,6 +280,7 @@ func (r *itemCollector) getResourceItems(log logrus.FieldLogger, gv schema.Group var items []*kubernetesResource for _, namespace := range namespacesToList { + // List items from Kubernetes API log = log.WithField("namespace", namespace) resourceClient, err := r.dynamicFactory.ClientForGroupVersionResource(gv, resource, namespace) @@ -287,18 +293,65 @@ func (r *itemCollector) getResourceItems(log logrus.FieldLogger, gv schema.Group if selector := r.backupRequest.Spec.LabelSelector; selector != nil { labelSelector = metav1.FormatLabelSelector(selector) } + listOptions := metav1.ListOptions{LabelSelector: labelSelector} log.Info("Listing items") - unstructuredList, err := resourceClient.List(metav1.ListOptions{LabelSelector: labelSelector}) - if err != nil { - log.WithError(errors.WithStack(err)).Error("Error listing items") - continue - } - log.Infof("Retrieved %d items", len(unstructuredList.Items)) + unstructuredItems := make([]unstructured.Unstructured, 0) - // collect the items - for i := range unstructuredList.Items { - item := &unstructuredList.Items[i] + if r.pageSize > 0 { + // If limit is positive, use a pager to split list over multiple requests + // Use Velero's dynamic list function instead of the default + listFunc := pager.SimplePageFunc(func(opts metav1.ListOptions) (runtime.Object, error) { + list, err := resourceClient.List(listOptions) + if err != nil { + return nil, err + } + return list, nil + }) + listPager := pager.New(listFunc) + // Use the page size defined in the server config + // TODO allow configuration of page buffer size + listPager.PageSize = int64(r.pageSize) + // Add each item to temporary slice + var items []unstructured.Unstructured + err := listPager.EachListItem(context.Background(), listOptions, func(object runtime.Object) error { + item, isUnstructured := object.(*unstructured.Unstructured) + if !isUnstructured { + // We should never hit this + log.Error("Got type other than Unstructured from pager func") + return nil + } + items = append(items, *item) + return nil + }) + if statusError, isStatusError := err.(*apierrors.StatusError); isStatusError && statusError.Status().Reason == metav1.StatusReasonExpired { + log.WithError(errors.WithStack(err)).Error("Error paging item list. Falling back on unpaginated list") + unstructuredList, err := resourceClient.List(listOptions) + if err != nil { + log.WithError(errors.WithStack(err)).Error("Error listing items") + continue + } + items = unstructuredList.Items + } else if err != nil { + log.WithError(errors.WithStack(err)).Error("Error paging item list") + continue + } + unstructuredItems = append(unstructuredItems, items...) + } else { + // If limit is not positive, do not use paging. Instead, request all items at once + unstructuredList, err := resourceClient.List(metav1.ListOptions{LabelSelector: labelSelector}) + unstructuredItems = append(unstructuredItems, unstructuredList.Items...) + if err != nil { + log.WithError(errors.WithStack(err)).Error("Error listing items") + continue + } + } + + log.Infof("Retrieved %d items", len(unstructuredItems)) + + // Collect items in included Namespaces + for i := range unstructuredItems { + item := &unstructuredItems[i] if gr == kuberesource.Namespaces && !r.backupRequest.NamespaceIncludesExcludes.ShouldInclude(item.GetName()) { log.WithField("name", item.GetName()).Info("Skipping namespace because it's excluded") diff --git a/pkg/cmd/server/server.go b/pkg/cmd/server/server.go index 7d1aec490..e0ac43a01 100644 --- a/pkg/cmd/server/server.go +++ b/pkg/cmd/server/server.go @@ -90,8 +90,9 @@ const ( defaultResourceTerminatingTimeout = 10 * time.Minute // server's client default qps and burst - defaultClientQPS float32 = 20.0 - defaultClientBurst int = 30 + defaultClientQPS float32 = 20.0 + defaultClientBurst int = 30 + defaultClientPageSize int = 500 defaultProfilerAddress = "localhost:6060" @@ -115,6 +116,7 @@ type serverConfig struct { disabledControllers []string clientQPS float32 clientBurst int + clientPageSize int profilerAddress string formatFlag *logging.FormatFlag defaultResticMaintenanceFrequency time.Duration @@ -142,6 +144,7 @@ func NewCommand(f client.Factory) *cobra.Command { restoreResourcePriorities: defaultRestorePriorities, clientQPS: defaultClientQPS, clientBurst: defaultClientBurst, + clientPageSize: defaultClientPageSize, profilerAddress: defaultProfilerAddress, resourceTerminatingTimeout: defaultResourceTerminatingTimeout, formatFlag: logging.NewFormatFlag(), @@ -205,6 +208,7 @@ func NewCommand(f client.Factory) *cobra.Command { command.Flags().Var(&volumeSnapshotLocations, "default-volume-snapshot-locations", "List of unique volume providers and default volume snapshot location (provider1:location-01,provider2:location-02,...)") command.Flags().Float32Var(&config.clientQPS, "client-qps", config.clientQPS, "Maximum number of requests per second by the server to the Kubernetes API once the burst limit has been reached.") command.Flags().IntVar(&config.clientBurst, "client-burst", config.clientBurst, "Maximum number of requests by the server to the Kubernetes API in a short period of time.") + command.Flags().IntVar(&config.clientPageSize, "client-page-size", config.clientPageSize, "Page size of requests by the server to the Kubernetes API when listing objects during a backup. Set to 0 to disable paging.") command.Flags().StringVar(&config.profilerAddress, "profiler-address", config.profilerAddress, "The address to expose the pprof profiler.") command.Flags().DurationVar(&config.resourceTerminatingTimeout, "terminating-resource-timeout", config.resourceTerminatingTimeout, "How long to wait on persistent volumes and namespaces to terminate during a restore before timing out.") command.Flags().DurationVar(&config.defaultBackupTTL, "default-backup-ttl", config.defaultBackupTTL, "How long to wait by default before backups can be garbage collected.") @@ -249,6 +253,10 @@ func newServer(f client.Factory, config serverConfig, logger *logrus.Logger) (*s } f.SetClientBurst(config.clientBurst) + if config.clientPageSize < 0 { + return nil, errors.New("client-page-size must not be negative") + } + kubeClient, err := f.KubeClient() if err != nil { return nil, err @@ -610,6 +618,7 @@ func (s *server) runControllers(defaultVolumeSnapshotLocations map[string]string s.resticManager, s.config.podVolumeOperationTimeout, s.config.defaultVolumesToRestic, + s.config.clientPageSize, ) cmd.CheckError(err) diff --git a/site/content/docs/main/backup-reference.md b/site/content/docs/main/backup-reference.md index 33572edcf..60e019188 100644 --- a/site/content/docs/main/backup-reference.md +++ b/site/content/docs/main/backup-reference.md @@ -60,3 +60,11 @@ velero backup create --from-schedule example-schedule ``` This command will immediately trigger a new backup based on your template for `example-schedule`. This will not affect the backup schedule, and another backup will trigger at the scheduled time. + +## Kubernetes API Pagination + +By default, Velero will paginate the LIST API call for each resource type in the Kubernetes API when collecting items into a backup. The `--client-page-size` flag for the Velero server configures the size of each page. + +Depending on the cluster's scale, tuning the page size can improve backup performance. You can experiment with higher values, noting their impact on the relevant `apiserver_request_duration_seconds_*` metrics from the Kubernetes apiserver. + +Pagination can be entirely disabled by setting `--client-page-size` to `0`. This will request all items in a single unpaginated LIST call.