Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions config/crd/bases/pgbackrest.cnpg.opera.com_archives.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -386,6 +386,18 @@ spec:
value - with 1 being the minimum accepted value.
minimum: 1
type: integer
maxPrefetch:
description: |-
Number of WAL files to download ahead of the one PostgreSQL requests while
restoring, when maxParallel is greater than 1. The downloads run in the
background and each one that completes starts the next, so one slow download
does not hold up the others; no more than maxParallel run at a time. The files
wait in the spool directory on the pod's ephemeral storage, so this also bounds
the space the spool takes, in WAL segments. When some downloads take much longer
than the rest, about four times maxParallel keeps all of them busy. Defaults to
maxParallel - 1.
minimum: 1
type: integer
restoreAdditionalCommandArgs:
description: |-
Additional arguments that can be appended to the 'pgbackrest restore'
Expand Down
131 changes: 69 additions & 62 deletions internal/cnpgi/common/wal.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ import (
pgbackrestv1 "github.com/operasoftware/cnpg-plugin-pgbackrest/api/v1"
"github.com/operasoftware/cnpg-plugin-pgbackrest/internal/cnpgi/metadata"
"github.com/operasoftware/cnpg-plugin-pgbackrest/internal/cnpgi/operator/config"
pgbackrestApi "github.com/operasoftware/cnpg-plugin-pgbackrest/internal/pgbackrest/api"
"github.com/operasoftware/cnpg-plugin-pgbackrest/internal/pgbackrest/archiver"
pgbackrestBackup "github.com/operasoftware/cnpg-plugin-pgbackrest/internal/pgbackrest/backup"
pgbackrestCommand "github.com/operasoftware/cnpg-plugin-pgbackrest/internal/pgbackrest/command"
Expand Down Expand Up @@ -306,9 +307,35 @@ func (w WALServiceImplementation) restoreFromPgbackrestArchive(
return fmt.Errorf("while creating the restorer: %w", err)
}

// Step 1: check if this WAL file is not already in the spool
// Step 1: keep the WAL files after this one downloading into the spool in the background,
// so that they are ready when PostgreSQL asks for them. This runs on every call, also when
// the file is in the spool, to move the prefetch forward as PostgreSQL does.
streamingAvailable := isStreamingAvailable(cluster, w.InstanceName)
prefetch, err := walsToPrefetch(walName, pgbackrestConfiguration.Wal, controlledPromotion)
if err != nil {
return fmt.Errorf("while generating the list of WAL files to restore: %w", err)
}
prefetch.Options = options
prefetch.OnDone = func(prefetchedWAL string, err error) {
// We skip this step if streaming connection is not available
if !streamingAvailable || !errors.Is(err, pgbackrestRestorer.ErrWALNotFound) {
return
}
contextLogger.Info(
"Set end-of-wal-stream flag as one of the WAL files to be prefetched was not found",
"walName", prefetchedWAL)
if err := walRestorer.SetEndOfWALStream(); err != nil {
contextLogger.Error(err, "while setting the end-of-wal-stream flag")
}
}
walRestorer.Prefetch(ctx, prefetch)

// Step 2: check if this WAL file is already in the spool, or being prefetched there
var wasInSpool bool
if wasInSpool, err = walRestorer.RestoreFromSpool(walName, destinationPath); err != nil {
if wasInSpool, err = walRestorer.RestoreFromSpool(ctx, walName, destinationPath); err != nil {
if errors.Is(err, pgbackrestRestorer.ErrWALNotFound) {
return newWALNotFoundError()
}
return fmt.Errorf("while restoring a file from the spool directory: %w", err)
}
if wasInSpool {
Expand All @@ -319,67 +346,29 @@ func (w WALServiceImplementation) restoreFromPgbackrestArchive(
}

// We skip this step if streaming connection is not available
if isStreamingAvailable(cluster, w.InstanceName) {
if streamingAvailable {
if err := checkEndOfWALStreamFlag(walRestorer); err != nil {
return err
}
}

// Step 3: gather the WAL files names to restore. If the required file isn't a regular WAL, we download it directly.
var walFilesList []string
maxParallel := 1
if pgbackrestConfiguration.Wal != nil && pgbackrestConfiguration.Wal.MaxParallel > 1 {
maxParallel = pgbackrestConfiguration.Wal.MaxParallel
}
if IsWALFile(walName) {
// If this is a regular WAL file, we try to prefetch
if walFilesList, err = gatherWALFilesToRestore(walName, maxParallel, controlledPromotion); err != nil {
return fmt.Errorf("while generating the list of WAL files to restore: %w", err)
}
} else {
// This is not a regular WAL file, we fetch it directly
walFilesList = []string{walName}
}

// Step 4: download the WAL files into the required place
// Step 3: restore the requested WAL file. PostgreSQL gets its answer as soon as the file
// is restored or known to be missing, while the prefetch goes on in the background.
downloadStartTime := time.Now()
walStatus := walRestorer.RestoreList(ctx, walFilesList, destinationPath, options)

// We return immediately if the first WAL has errors, because the first WAL
// is the one that PostgreSQL has requested to restore.
// The failure has already been logged in walRestorer.RestoreList method
if walStatus[0].Err != nil {
if errors.Is(walStatus[0].Err, pgbackrestRestorer.ErrWALNotFound) {
prefetching, err := walRestorer.RestoreWithPrefetch(ctx, walName, destinationPath, prefetch)
if err != nil {
// The failure has already been logged in walRestorer.RestoreWithPrefetch method
if errors.Is(err, pgbackrestRestorer.ErrWALNotFound) {
return newWALNotFoundError()
}

return walStatus[0].Err
}

// We skip this step if streaming connection is not available
endOfWALStream := isEndOfWALStream(walStatus)
if isStreamingAvailable(cluster, w.InstanceName) && endOfWALStream {
contextLogger.Info(
"Set end-of-wal-stream flag as one of the WAL files to be prefetched was not found")

err = walRestorer.SetEndOfWALStream()
if err != nil {
return err
}
}

successfulWalRestore := 0
for idx := range walStatus {
if walStatus[idx].Err == nil {
successfulWalRestore++
}
return err
}

contextLogger.Info("WAL restore command completed (parallel)",
"walName", walName,
"maxParallel", maxParallel,
"successfulWalRestore", successfulWalRestore,
"failedWalRestore", maxParallel-successfulWalRestore,
"maxParallel", prefetch.Parallel,
"prefetching", prefetching,
"startTime", startTime,
"downloadStartTime", downloadStartTime,
"downloadTotalTime", time.Since(downloadStartTime),
Expand Down Expand Up @@ -477,6 +466,36 @@ func isStreamingAvailable(cluster *cnpgv1.Cluster, podName string) bool {
return false
}

// walsToPrefetch returns what to prefetch into the spool for a request of walName: the next
// maxPrefetch WAL files, maxParallel - 1 when it is not set, downloaded up to maxParallel at
// a time. With maxParallel 1, or for a file that is not a regular WAL, nothing is prefetched
// but the partial WAL of a controlled promotion.
func walsToPrefetch(
walName string,
configuration *pgbackrestApi.WalBackupConfiguration,
controlledPromotion bool,
) (pgbackrestRestorer.Prefetch, error) {
prefetch := pgbackrestRestorer.Prefetch{Parallel: 1}
lookahead := 0
if configuration != nil && configuration.MaxParallel > 1 {
prefetch.Parallel = configuration.MaxParallel
lookahead = configuration.MaxParallel - 1
if configuration.MaxPrefetch > 0 {
lookahead = configuration.MaxPrefetch
}
}
if !IsWALFile(walName) {
return prefetch, nil
}

walList, err := gatherWALFilesToRestore(walName, lookahead+1, controlledPromotion)
if err != nil {
return prefetch, err
}
prefetch.WALs = walList[1:]
return prefetch, nil
}

// gatherWALFilesToRestore files a list of possible WAL files to restore, always
// including as the first one the requested WAL file.
func gatherWALFilesToRestore(walName string, parallel int, controlledPromotion bool) (walList []string, err error) {
Expand Down Expand Up @@ -526,15 +545,3 @@ func checkEndOfWALStreamFlag(walRestorer *pgbackrestRestorer.WALRestorer) error
}
return nil
}

// isEndOfWALStream returns true if one of the downloads has returned
// a file-not-found error.
func isEndOfWALStream(results []pgbackrestRestorer.Result) bool {
for _, result := range results {
if errors.Is(result.Err, pgbackrestRestorer.ErrWALNotFound) {
return true
}
}

return false
}
12 changes: 12 additions & 0 deletions internal/pgbackrest/api/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,18 @@ type WalBackupConfiguration struct {
// +kubebuilder:validation:Minimum=1
// +optional
MaxParallel int `json:"maxParallel,omitempty"`

// Number of WAL files to download ahead of the one PostgreSQL requests while
// restoring, when maxParallel is greater than 1. The downloads run in the
// background and each one that completes starts the next, so one slow download
// does not hold up the others; no more than maxParallel run at a time. The files
// wait in the spool directory on the pod's ephemeral storage, so this also bounds
// the space the spool takes, in WAL segments. When some downloads take much longer
// than the rest, about four times maxParallel keeps all of them busy. Defaults to
// maxParallel - 1.
// +kubebuilder:validation:Minimum=1
// +optional
MaxPrefetch int `json:"maxPrefetch,omitempty"`
// Additional arguments that can be appended to the 'pgbackrest archive-push'
// command-line invocation. These arguments provide flexibility to customize
// the WAL archive process further, according to specific requirements or configurations.
Expand Down
Loading
Loading