mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-08 08:22:30 +02:00
Add logical backup (#442)
* Add k8s cron job to spawn logical backups * Minor doc updates
This commit is contained in:
+56
-2
@@ -81,6 +81,7 @@ type Cluster struct {
|
||||
currentProcess Process
|
||||
processMu sync.RWMutex // protects the current operation for reporting, no need to hold the master mutex
|
||||
specMu sync.RWMutex // protects the spec for reporting, no need to hold the master mutex
|
||||
|
||||
}
|
||||
|
||||
type compareStatefulsetResult struct {
|
||||
@@ -298,6 +299,13 @@ func (c *Cluster) Create() error {
|
||||
c.logger.Infof("databases have been successfully created")
|
||||
}
|
||||
|
||||
if c.Postgresql.Spec.EnableLogicalBackup {
|
||||
if err := c.createLogicalBackupJob(); err != nil {
|
||||
return fmt.Errorf("could not create a k8s cron job for logical backups: %v", err)
|
||||
}
|
||||
c.logger.Info("a k8s cron job for logical backup has been successfully created")
|
||||
}
|
||||
|
||||
if err := c.listResources(); err != nil {
|
||||
c.logger.Errorf("could not list resources: %v", err)
|
||||
}
|
||||
@@ -481,8 +489,10 @@ func compareResoucesAssumeFirstNotNil(a *v1.ResourceRequirements, b *v1.Resource
|
||||
|
||||
}
|
||||
|
||||
// Update changes Kubernetes objects according to the new specification. Unlike the sync case, the missing object.
|
||||
// (i.e. service) is treated as an error.
|
||||
// Update changes Kubernetes objects according to the new specification. Unlike the sync case, the missing object
|
||||
// (i.e. service) is treated as an error
|
||||
// logical backup cron jobs are an exception: a user-initiated Update can enable a logical backup job
|
||||
// for a cluster that had no such job before. In this case a missing job is not an error.
|
||||
func (c *Cluster) Update(oldSpec, newSpec *acidv1.Postgresql) error {
|
||||
updateFailed := false
|
||||
|
||||
@@ -569,6 +579,43 @@ func (c *Cluster) Update(oldSpec, newSpec *acidv1.Postgresql) error {
|
||||
}
|
||||
}()
|
||||
|
||||
// logical backup job
|
||||
func() {
|
||||
|
||||
// create if it did not exist
|
||||
if !oldSpec.Spec.EnableLogicalBackup && newSpec.Spec.EnableLogicalBackup {
|
||||
c.logger.Debugf("creating backup cron job")
|
||||
if err := c.createLogicalBackupJob(); err != nil {
|
||||
c.logger.Errorf("could not create a k8s cron job for logical backups: %v", err)
|
||||
updateFailed = true
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// delete if no longer needed
|
||||
if oldSpec.Spec.EnableLogicalBackup && !newSpec.Spec.EnableLogicalBackup {
|
||||
c.logger.Debugf("deleting backup cron job")
|
||||
if err := c.deleteLogicalBackupJob(); err != nil {
|
||||
c.logger.Errorf("could not delete a k8s cron job for logical backups: %v", err)
|
||||
updateFailed = true
|
||||
return
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
// apply schedule changes
|
||||
// this is the only parameter of logical backups a user can overwrite in the cluster manifest
|
||||
if (oldSpec.Spec.EnableLogicalBackup && newSpec.Spec.EnableLogicalBackup) &&
|
||||
(newSpec.Spec.LogicalBackupSchedule != oldSpec.Spec.LogicalBackupSchedule) {
|
||||
c.logger.Debugf("updating schedule of the backup cron job")
|
||||
if err := c.syncLogicalBackupJob(); err != nil {
|
||||
c.logger.Errorf("could not sync logical backup jobs: %v", err)
|
||||
updateFailed = true
|
||||
}
|
||||
}
|
||||
|
||||
}()
|
||||
|
||||
// Roles and Databases
|
||||
if !(c.databaseAccessDisabled() || c.getNumberOfInstances(&c.Spec) <= 0) {
|
||||
c.logger.Debugf("syncing roles")
|
||||
@@ -597,6 +644,12 @@ func (c *Cluster) Delete() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
// delete the backup job before the stateful set of the cluster to prevent connections to non-existing pods
|
||||
// deleting the cron job also removes pods and batch jobs it created
|
||||
if err := c.deleteLogicalBackupJob(); err != nil {
|
||||
c.logger.Warningf("could not remove the logical backup k8s cron job; %v", err)
|
||||
}
|
||||
|
||||
if err := c.deleteStatefulSet(); err != nil {
|
||||
c.logger.Warningf("could not delete statefulset: %v", err)
|
||||
}
|
||||
@@ -629,6 +682,7 @@ func (c *Cluster) Delete() {
|
||||
if err := c.deletePatroniClusterObjects(); err != nil {
|
||||
c.logger.Warningf("could not remove leftover patroni objects; %v", err)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
//NeedsRepair returns true if the cluster should be included in the repair scan (based on its in-memory status).
|
||||
|
||||
+168
-2
@@ -20,6 +20,8 @@ import (
|
||||
"github.com/zalando/postgres-operator/pkg/util"
|
||||
"github.com/zalando/postgres-operator/pkg/util/config"
|
||||
"github.com/zalando/postgres-operator/pkg/util/constants"
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
batchv1beta1 "k8s.io/api/batch/v1beta1"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
)
|
||||
|
||||
@@ -352,7 +354,7 @@ func generateVolumeMounts() []v1.VolumeMount {
|
||||
}
|
||||
}
|
||||
|
||||
func generateSpiloContainer(
|
||||
func generateContainer(
|
||||
name string,
|
||||
dockerImage *string,
|
||||
resourceRequirements *v1.ResourceRequirements,
|
||||
@@ -792,7 +794,7 @@ func (c *Cluster) generateStatefulSet(spec *acidv1.PostgresSpec) (*v1beta1.State
|
||||
|
||||
// generate the spilo container
|
||||
c.logger.Debugf("Generating Spilo container, environment variables: %v", spiloEnvVars)
|
||||
spiloContainer := generateSpiloContainer(c.containerName(),
|
||||
spiloContainer := generateContainer(c.containerName(),
|
||||
&effectiveDockerImage,
|
||||
resourceRequirements,
|
||||
spiloEnvVars,
|
||||
@@ -1281,3 +1283,167 @@ func (c *Cluster) getClusterServiceConnectionParameters(clusterName string) (hos
|
||||
port = "5432"
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) generateLogicalBackupJob() (*batchv1beta1.CronJob, error) {
|
||||
|
||||
var (
|
||||
err error
|
||||
podTemplate *v1.PodTemplateSpec
|
||||
resourceRequirements *v1.ResourceRequirements
|
||||
)
|
||||
|
||||
// NB: a cron job creates standard batch jobs according to schedule; these batch jobs manage pods and clean-up
|
||||
|
||||
c.logger.Debug("Generating logical backup pod template")
|
||||
|
||||
// allocate for the backup pod the same amount of resources as for normal DB pods
|
||||
defaultResources := c.makeDefaultResources()
|
||||
resourceRequirements, err = generateResourceRequirements(c.Spec.Resources, defaultResources)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("could not generate resource requirements for logical backup pods: %v", err)
|
||||
}
|
||||
|
||||
envVars := c.generateLogicalBackupPodEnvVars()
|
||||
logicalBackupContainer := generateContainer(
|
||||
"logical-backup",
|
||||
&c.OpConfig.LogicalBackup.LogicalBackupDockerImage,
|
||||
resourceRequirements,
|
||||
envVars,
|
||||
[]v1.VolumeMount{},
|
||||
c.OpConfig.SpiloPrivileged, // use same value as for normal DB pods
|
||||
)
|
||||
|
||||
labels := map[string]string{
|
||||
"version": c.Name,
|
||||
"application": "spilo-logical-backup",
|
||||
}
|
||||
podAffinityTerm := v1.PodAffinityTerm{
|
||||
LabelSelector: &metav1.LabelSelector{
|
||||
MatchLabels: labels,
|
||||
},
|
||||
TopologyKey: "kubernetes.io/hostname",
|
||||
}
|
||||
podAffinity := v1.Affinity{
|
||||
PodAffinity: &v1.PodAffinity{
|
||||
PreferredDuringSchedulingIgnoredDuringExecution: []v1.WeightedPodAffinityTerm{{
|
||||
Weight: 1,
|
||||
PodAffinityTerm: podAffinityTerm,
|
||||
},
|
||||
},
|
||||
}}
|
||||
|
||||
// re-use the method that generates DB pod templates
|
||||
if podTemplate, err = generatePodTemplate(
|
||||
c.Namespace,
|
||||
c.labelsSet(true),
|
||||
logicalBackupContainer,
|
||||
[]v1.Container{},
|
||||
[]v1.Container{},
|
||||
&[]v1.Toleration{},
|
||||
nodeAffinity(c.OpConfig.NodeReadinessLabel),
|
||||
int64(c.OpConfig.PodTerminateGracePeriod.Seconds()),
|
||||
c.OpConfig.PodServiceAccountName,
|
||||
c.OpConfig.KubeIAMRole,
|
||||
"",
|
||||
false,
|
||||
false,
|
||||
""); err != nil {
|
||||
return nil, fmt.Errorf("could not generate pod template for logical backup pod: %v", err)
|
||||
}
|
||||
|
||||
// overwrite specifc params of logical backups pods
|
||||
podTemplate.Spec.Affinity = &podAffinity
|
||||
podTemplate.Spec.RestartPolicy = "Never" // affects containers within a pod
|
||||
|
||||
// configure a batch job
|
||||
|
||||
jobSpec := batchv1.JobSpec{
|
||||
Template: *podTemplate,
|
||||
}
|
||||
|
||||
// configure a cron job
|
||||
|
||||
jobTemplateSpec := batchv1beta1.JobTemplateSpec{
|
||||
Spec: jobSpec,
|
||||
}
|
||||
|
||||
schedule := c.Postgresql.Spec.LogicalBackupSchedule
|
||||
if schedule == "" {
|
||||
schedule = c.OpConfig.LogicalBackupSchedule
|
||||
}
|
||||
|
||||
cronJob := &batchv1beta1.CronJob{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
Name: c.getLogicalBackupJobName(),
|
||||
Namespace: c.Namespace,
|
||||
Labels: c.labelsSet(true),
|
||||
},
|
||||
Spec: batchv1beta1.CronJobSpec{
|
||||
Schedule: schedule,
|
||||
JobTemplate: jobTemplateSpec,
|
||||
ConcurrencyPolicy: batchv1beta1.ForbidConcurrent,
|
||||
},
|
||||
}
|
||||
|
||||
return cronJob, nil
|
||||
}
|
||||
|
||||
func (c *Cluster) generateLogicalBackupPodEnvVars() []v1.EnvVar {
|
||||
|
||||
envVars := []v1.EnvVar{
|
||||
{
|
||||
Name: "SCOPE",
|
||||
Value: c.Name,
|
||||
},
|
||||
// Bucket env vars
|
||||
{
|
||||
Name: "LOGICAL_BACKUP_S3_BUCKET",
|
||||
Value: c.OpConfig.LogicalBackup.LogicalBackupS3Bucket,
|
||||
},
|
||||
{
|
||||
Name: "LOGICAL_BACKUP_S3_BUCKET_SCOPE_SUFFIX",
|
||||
Value: getBucketScopeSuffix(string(c.Postgresql.GetUID())),
|
||||
},
|
||||
// Postgres env vars
|
||||
{
|
||||
Name: "PG_VERSION",
|
||||
Value: c.Spec.PgVersion,
|
||||
},
|
||||
{
|
||||
Name: "PGPORT",
|
||||
Value: "5432",
|
||||
},
|
||||
{
|
||||
Name: "PGUSER",
|
||||
Value: c.OpConfig.SuperUsername,
|
||||
},
|
||||
{
|
||||
Name: "PGDATABASE",
|
||||
Value: c.OpConfig.SuperUsername,
|
||||
},
|
||||
{
|
||||
Name: "PGSSLMODE",
|
||||
Value: "require",
|
||||
},
|
||||
{
|
||||
Name: "PGPASSWORD",
|
||||
ValueFrom: &v1.EnvVarSource{
|
||||
SecretKeyRef: &v1.SecretKeySelector{
|
||||
LocalObjectReference: v1.LocalObjectReference{
|
||||
Name: c.credentialSecretName(c.OpConfig.SuperUsername),
|
||||
},
|
||||
Key: "password",
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
c.logger.Debugf("Generated logical backup env vars %v", envVars)
|
||||
|
||||
return envVars
|
||||
}
|
||||
|
||||
// getLogicalBackupJobName returns the name; the job itself may not exists
|
||||
func (c *Cluster) getLogicalBackupJobName() (jobName string) {
|
||||
return "logical-backup-" + c.clusterName().Name
|
||||
}
|
||||
|
||||
@@ -6,7 +6,8 @@ import (
|
||||
"strings"
|
||||
|
||||
"k8s.io/api/apps/v1beta1"
|
||||
"k8s.io/api/core/v1"
|
||||
batchv1beta1 "k8s.io/api/batch/v1beta1"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
@@ -609,6 +610,51 @@ func (c *Cluster) createRoles() (err error) {
|
||||
return c.syncRoles()
|
||||
}
|
||||
|
||||
func (c *Cluster) createLogicalBackupJob() (err error) {
|
||||
|
||||
c.setProcessName("creating a k8s cron job for logical backups")
|
||||
|
||||
logicalBackupJobSpec, err := c.generateLogicalBackupJob()
|
||||
if err != nil {
|
||||
return fmt.Errorf("could not generate k8s cron job spec: %v", err)
|
||||
}
|
||||
c.logger.Debugf("Generated cronJobSpec: %v", logicalBackupJobSpec)
|
||||
|
||||
_, err = c.KubeClient.CronJobsGetter.CronJobs(c.Namespace).Create(logicalBackupJobSpec)
|
||||
if err != nil {
|
||||
return fmt.Errorf("could not create k8s cron job: %v", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) patchLogicalBackupJob(newJob *batchv1beta1.CronJob) error {
|
||||
c.setProcessName("patching logical backup job")
|
||||
|
||||
patchData, err := specPatch(newJob.Spec)
|
||||
if err != nil {
|
||||
return fmt.Errorf("could not form patch for the logical backup job: %v", err)
|
||||
}
|
||||
|
||||
// update the backup job spec
|
||||
_, err = c.KubeClient.CronJobsGetter.CronJobs(c.Namespace).Patch(
|
||||
c.getLogicalBackupJobName(),
|
||||
types.MergePatchType,
|
||||
patchData, "")
|
||||
if err != nil {
|
||||
return fmt.Errorf("could not patch logical backup job: %v", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) deleteLogicalBackupJob() error {
|
||||
|
||||
c.logger.Info("removing the logical backup job")
|
||||
|
||||
return c.KubeClient.CronJobsGetter.CronJobs(c.Namespace).Delete(c.getLogicalBackupJobName(), c.deleteOptions)
|
||||
}
|
||||
|
||||
// GetServiceMaster returns cluster's kubernetes master Service
|
||||
func (c *Cluster) GetServiceMaster() *v1.Service {
|
||||
return c.Services[Master]
|
||||
|
||||
+65
-1
@@ -3,7 +3,8 @@ package cluster
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"k8s.io/api/core/v1"
|
||||
batchv1beta1 "k8s.io/api/batch/v1beta1"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
|
||||
@@ -92,6 +93,16 @@ func (c *Cluster) Sync(newSpec *acidv1.Postgresql) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// create a logical backup job unless we are running without pods or disable that feature explicitly
|
||||
if c.Spec.EnableLogicalBackup && c.getNumberOfInstances(&c.Spec) > 0 {
|
||||
|
||||
c.logger.Debug("syncing logical backup job")
|
||||
if err = c.syncLogicalBackupJob(); err != nil {
|
||||
err = fmt.Errorf("could not sync the logical backup job: %v", err)
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -519,3 +530,56 @@ func (c *Cluster) syncDatabases() error {
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) syncLogicalBackupJob() error {
|
||||
var (
|
||||
job *batchv1beta1.CronJob
|
||||
desiredJob *batchv1beta1.CronJob
|
||||
err error
|
||||
)
|
||||
c.setProcessName("syncing the logical backup job")
|
||||
|
||||
// sync the job if it exists
|
||||
|
||||
jobName := c.getLogicalBackupJobName()
|
||||
if job, err = c.KubeClient.CronJobsGetter.CronJobs(c.Namespace).Get(jobName, metav1.GetOptions{}); err == nil {
|
||||
|
||||
desiredJob, err = c.generateLogicalBackupJob()
|
||||
if err != nil {
|
||||
return fmt.Errorf("could not generate the desired logical backup job state: %v", err)
|
||||
}
|
||||
if match, reason := k8sutil.SameLogicalBackupJob(job, desiredJob); !match {
|
||||
c.logger.Infof("logical job %q is not in the desired state and needs to be updated",
|
||||
c.getLogicalBackupJobName(),
|
||||
)
|
||||
if reason != "" {
|
||||
c.logger.Infof("reason: %s", reason)
|
||||
}
|
||||
if err = c.patchLogicalBackupJob(desiredJob); err != nil {
|
||||
return fmt.Errorf("could not update logical backup job to match desired state: %v", err)
|
||||
}
|
||||
c.logger.Info("the logical backup job is synced")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if !k8sutil.ResourceNotFound(err) {
|
||||
return fmt.Errorf("could not get logical backp job: %v", err)
|
||||
}
|
||||
|
||||
// no existing logical backup job, create new one
|
||||
c.logger.Info("could not find the cluster's logical backup job")
|
||||
|
||||
if err = c.createLogicalBackupJob(); err == nil {
|
||||
c.logger.Infof("created missing logical backup job %q", jobName)
|
||||
} else {
|
||||
if !k8sutil.ResourceAlreadyExists(err) {
|
||||
return fmt.Errorf("could not create missing logical backup job: %v", err)
|
||||
}
|
||||
c.logger.Infof("logical backup job %q already exists", jobName)
|
||||
if _, err = c.KubeClient.CronJobsGetter.CronJobs(c.Namespace).Get(jobName, metav1.GetOptions{}); err != nil {
|
||||
return fmt.Errorf("could not fetch existing logical backup job: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
+1
-1
@@ -12,7 +12,7 @@ import (
|
||||
"time"
|
||||
|
||||
"k8s.io/api/apps/v1beta1"
|
||||
"k8s.io/api/core/v1"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
policybeta1 "k8s.io/api/policy/v1beta1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
|
||||
Reference in New Issue
Block a user