Skip to content

Commit 70de743

Browse files
grzegorz8Grzegorz Kołakowski
andauthored
Use stop-with-savepoint for automatic FlinkCluster updates (#1042)
Replace the automatic update path's TriggerSavepoint(..., cancel=true) call with StopJobWithSavepoint(). The /stop endpoint coordinates the final savepoint with job termination, preventing records from being processed after the savepoint barrier. Co-authored-by: Grzegorz Kołakowski <grzegorzk@spotify.com>
1 parent 6fadd2c commit 70de743

3 files changed

Lines changed: 328 additions & 23 deletions

File tree

controllers/flinkcluster/flinkcluster_reconciler.go

Lines changed: 35 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -512,18 +512,24 @@ func (reconciler *ClusterReconciler) reconcileJob(ctx context.Context) (ctrl.Res
512512
observedSubmitter := observed.flinkJobSubmitter.job
513513

514514
if desiredJob != nil && job.IsTerminated(jobSpec) {
515-
// When the job was cancelled as part of an update (savepoint trigger
516-
// reason is "update"), don't treat it as terminally done. The operator
517-
// needs to proceed with creating a new job submitter to restart the
518-
// job from the savepoint taken during the update.
515+
// When the job was stopped as part of an update (savepoint trigger
516+
// reason is "update"), don't treat it as terminally done. Flink reports
517+
// a job stopped through /stop as FINISHED, which maps to Succeeded.
518+
// The operator needs to proceed with creating a new job submitter to
519+
// restart the job from the savepoint taken during the update.
520+
// Stop-with-savepoint normally produces Succeeded. Accept Cancelled as well
521+
// for updates started by older operator versions and unexpected cancellation
522+
// after the final update savepoint completed.
519523
savepointStatus := recorded.Savepoint
520-
cancelledForUpdate := job.State == v1beta1.JobStateCancelled &&
524+
stoppedForUpdate := job.FinalSavepoint &&
525+
(job.State == v1beta1.JobStateCancelled ||
526+
job.State == v1beta1.JobStateSucceeded) &&
521527
savepointStatus != nil &&
522528
savepointStatus.TriggerReason == v1beta1.SavepointReasonUpdate
523-
if !cancelledForUpdate {
529+
if !stoppedForUpdate {
524530
return ctrl.Result{}, nil
525531
}
526-
log.Info("Job was cancelled for update, proceeding with job resubmission")
532+
log.Info("Job was stopped for update, proceeding with job resubmission")
527533
}
528534

529535
if wasJobCancelRequested(observed.cluster.Status.Control) {
@@ -726,12 +732,12 @@ func (reconciler *ClusterReconciler) trySuspendJob(ctx context.Context) (*v1beta
726732
log.Info("Checking the conditions for progressing")
727733
var canSuspend = reconciler.canSuspendJob(ctx, jobID, recorded.Savepoint)
728734
if canSuspend {
729-
log.Info("Triggering savepoint for suspending job")
735+
log.Info("Stopping job with savepoint for update")
730736
var newSavepointStatus, err = reconciler.triggerSavepoint(ctx, jobID, v1beta1.SavepointReasonUpdate, true)
731737
if err != nil {
732-
log.Info("Failed to trigger savepoint", "jobID", jobID, "triggerID", newSavepointStatus.TriggerID, "error", err)
738+
log.Info("Failed to stop job with savepoint", "jobID", jobID, "triggerID", newSavepointStatus.TriggerID, "error", err)
733739
} else {
734-
log.Info("Successfully savepoint triggered", "jobID", jobID, "triggerID", newSavepointStatus.TriggerID)
740+
log.Info("Successfully triggered stop with savepoint", "jobID", jobID, "triggerID", newSavepointStatus.TriggerID)
735741
}
736742
return newSavepointStatus, err
737743
}
@@ -806,16 +812,13 @@ func (reconciler *ClusterReconciler) cancelFlinkJob(ctx context.Context, jobID s
806812

807813
if takeSavepoint && canTakeSavepoint(reconciler.observed.cluster) {
808814
log.Info("Stopping job with savepoint", "jobID", jobID)
809-
formatType := savepointFormatType(reconciler.observed.cluster)
810-
triggerID, err := reconciler.flinkClient.StopJobWithSavepoint(
811-
apiBaseURL, jobID, *reconciler.observed.cluster.Spec.Job.SavepointsDir, string(formatType))
815+
newSavepointStatus, err := reconciler.triggerSavepoint(ctx, jobID, v1beta1.SavepointReasonJobCancel, true)
812816
if err != nil {
813817
return err
814818
}
815-
newSavepointStatus := reconciler.getNewSavepointStatus(triggerID.RequestID, v1beta1.SavepointReasonJobCancel, "", true, formatType)
816819
var newControlStatus *v1beta1.FlinkClusterControlStatus
817820
reconciler.updateStatus(ctx, &newSavepointStatus, &newControlStatus)
818-
location, err := reconciler.waitForSavepointCompleted(ctx, apiBaseURL, jobID, triggerID.RequestID)
821+
location, err := reconciler.waitForSavepointCompleted(ctx, apiBaseURL, jobID, newSavepointStatus.TriggerID)
819822
reconciler.updateFinalSavepointStatus(ctx, newSavepointStatus, location, err)
820823
return err
821824
}
@@ -873,7 +876,7 @@ func (reconciler *ClusterReconciler) waitForSavepointCompleted(ctx context.Conte
873876
}
874877
}
875878

876-
// canSuspendJob
879+
// canSuspendJob reports whether a new stop-with-savepoint request should be triggered.
877880
func (reconciler *ClusterReconciler) canSuspendJob(ctx context.Context, jobID string, s *v1beta1.SavepointStatus) bool {
878881
log := logr.FromContextOrDiscard(ctx)
879882
var firstTry = !finalSavepointRequested(jobID, s)
@@ -883,8 +886,13 @@ func (reconciler *ClusterReconciler) canSuspendJob(ctx context.Context, jobID st
883886

884887
switch s.State {
885888
case v1beta1.SavepointStateSucceeded:
889+
job := reconciler.observed.cluster.Status.Components.Job
890+
if job != nil && !job.FinalSavepoint {
891+
log.Info("Previous update savepoint does not belong to the current job execution")
892+
return true
893+
}
886894
log.Info("Successfully savepoint completed, wait until the job stops")
887-
return true
895+
return false
888896
case v1beta1.SavepointStateInProgress:
889897
log.Info("Savepoint is in progress, wait until it is completed")
890898
return false
@@ -955,12 +963,13 @@ func (reconciler *ClusterReconciler) shouldTakeSavepoint() v1beta1.SavepointReas
955963
return ""
956964
}
957965

958-
// Trigger savepoint for a job then return savepoint status to update.
966+
// Trigger a savepoint for a job, using /stop when stop is true, then return
967+
// the savepoint status to update.
959968
func (reconciler *ClusterReconciler) triggerSavepoint(
960969
ctx context.Context,
961970
jobID string,
962971
triggerReason v1beta1.SavepointReason,
963-
cancel bool) (*v1beta1.SavepointStatus, error) {
972+
stop bool) (*v1beta1.SavepointStatus, error) {
964973
log := logr.FromContextOrDiscard(ctx)
965974
var cluster = reconciler.observed.cluster
966975
var apiBaseURL = getFlinkAPIBaseURL(reconciler.observed.cluster)
@@ -971,7 +980,13 @@ func (reconciler *ClusterReconciler) triggerSavepoint(
971980
var formatType = savepointFormatType(cluster)
972981
var err error
973982
log.Info(fmt.Sprintf("Trigger savepoint for %s", triggerReason), "jobID", jobID)
974-
savepointTriggerID, err = reconciler.flinkClient.TriggerSavepoint(apiBaseURL, jobID, *cluster.Spec.Job.SavepointsDir, cancel, string(formatType))
983+
if stop {
984+
savepointTriggerID, err = reconciler.flinkClient.StopJobWithSavepoint(
985+
apiBaseURL, jobID, *cluster.Spec.Job.SavepointsDir, string(formatType))
986+
} else {
987+
savepointTriggerID, err = reconciler.flinkClient.TriggerSavepoint(
988+
apiBaseURL, jobID, *cluster.Spec.Job.SavepointsDir, false, string(formatType))
989+
}
975990
if err != nil {
976991
// limit message size to 1KiB
977992
if message = err.Error(); len(message) > 1024 {

0 commit comments

Comments
 (0)