diff --git a/deploy/crd.yaml b/deploy/crd.yaml index 2d05d90d..cae2170f 100644 --- a/deploy/crd.yaml +++ b/deploy/crd.yaml @@ -78,6 +78,9 @@ spec: parallelism: type: integer minimum: 1 + updateMode: + type: string + enum: [Savepoint, Checkpoint, NoStateRestore] deleteMode: type: string enum: [Savepoint, None, ForceCancel] diff --git a/pkg/apis/app/v1beta1/types.go b/pkg/apis/app/v1beta1/types.go index 8b646b1b..9412131b 100644 --- a/pkg/apis/app/v1beta1/types.go +++ b/pkg/apis/app/v1beta1/types.go @@ -54,6 +54,7 @@ type FlinkApplicationSpec struct { Volumes []apiv1.Volume `json:"volumes,omitempty"` VolumeMounts []apiv1.VolumeMount `json:"volumeMounts,omitempty"` RestartNonce string `json:"restartNonce"` + UpdateMode UpdateMode `json:"updateMode,omitempty"` DeleteMode DeleteMode `json:"deleteMode,omitempty"` AllowNonRestoredState bool `json:"allowNonRestoredState,omitempty"` ForceRollback bool `json:"forceRollback"` @@ -61,6 +62,13 @@ type FlinkApplicationSpec struct { TearDownVersionHash string `json:"tearDownVersionHash,omitempty"` } +func (spec FlinkApplicationSpec) SavepointingDisabled() bool { + if spec.SavepointDisabled { + return spec.SavepointDisabled + } + return spec.UpdateMode == UpdateModeNoState +} + type FlinkConfig map[string]interface{} // Workaround for https://github.com/kubernetes-sigs/kubebuilder/issues/528 @@ -302,6 +310,14 @@ const ( DeleteModeNone DeleteMode = "None" ) +type UpdateMode string + +const ( + UpdateModeSavepoint UpdateMode = "Savepoint" + UpdateModeCheckpoint UpdateMode = "Checkpoint" + UpdateModeNoState UpdateMode = "NoStateRestore" +) + type HealthStatus string const ( @@ -353,6 +369,7 @@ const ( GetJobConfig FlinkMethod = "GetJobConfig" GetTaskManagers FlinkMethod = "GetTaskManagers" GetCheckpointCounts FlinkMethod = "GetCheckpointCounts" + GetCheckpointConfig FlinkMethod = "GetCheckpointConfig" GetJobOverview FlinkMethod = "GetJobOverview" SavepointJob FlinkMethod = "SavepointJob" ) diff --git a/pkg/controller/flink/client/api.go b/pkg/controller/flink/client/api.go index 5084f768..ff047632 100644 --- a/pkg/controller/flink/client/api.go +++ b/pkg/controller/flink/client/api.go @@ -58,6 +58,7 @@ type FlinkAPIInterface interface { GetJobConfig(ctx context.Context, url string, jobID string) (*JobConfigResponse, error) GetTaskManagers(ctx context.Context, url string) (*TaskManagersResponse, error) GetCheckpointCounts(ctx context.Context, url string, jobID string) (*CheckpointResponse, error) + GetCheckpointConfig(ctx context.Context, url string, jobID string) (*CheckpointConfigResponse, error) GetJobOverview(ctx context.Context, url string, jobID string) (*FlinkJobOverview, error) } @@ -66,49 +67,53 @@ type FlinkJobManagerClient struct { } type flinkJobManagerClientMetrics struct { - scope promutils.Scope - submitJobSuccessCounter labeled.Counter - submitJobFailureCounter labeled.Counter - cancelJobSuccessCounter labeled.Counter - cancelJobFailureCounter labeled.Counter - forceCancelJobSuccessCounter labeled.Counter - forceCancelJobFailureCounter labeled.Counter - checkSavepointSuccessCounter labeled.Counter - checkSavepointFailureCounter labeled.Counter - getJobsSuccessCounter labeled.Counter - getJobsFailureCounter labeled.Counter - getJobConfigSuccessCounter labeled.Counter - getJobConfigFailureCounter labeled.Counter - getClusterSuccessCounter labeled.Counter - getClusterFailureCounter labeled.Counter - getCheckpointsSuccessCounter labeled.Counter - getCheckpointsFailureCounter labeled.Counter - savepointJobSuccessCounter labeled.Counter - savepointJobFailureCounter labeled.Counter + scope promutils.Scope + submitJobSuccessCounter labeled.Counter + submitJobFailureCounter labeled.Counter + cancelJobSuccessCounter labeled.Counter + cancelJobFailureCounter labeled.Counter + forceCancelJobSuccessCounter labeled.Counter + forceCancelJobFailureCounter labeled.Counter + checkSavepointSuccessCounter labeled.Counter + checkSavepointFailureCounter labeled.Counter + getJobsSuccessCounter labeled.Counter + getJobsFailureCounter labeled.Counter + getJobConfigSuccessCounter labeled.Counter + getJobConfigFailureCounter labeled.Counter + getClusterSuccessCounter labeled.Counter + getClusterFailureCounter labeled.Counter + getCheckpointsSuccessCounter labeled.Counter + getCheckpointsFailureCounter labeled.Counter + savepointJobSuccessCounter labeled.Counter + savepointJobFailureCounter labeled.Counter + getCheckpointsConfigSuccessCounter labeled.Counter + getCheckpointsConfigFailureCounter labeled.Counter } func newFlinkJobManagerClientMetrics(scope promutils.Scope) *flinkJobManagerClientMetrics { flinkJmClientScope := scope.NewSubScope("flink_jm_client") return &flinkJobManagerClientMetrics{ - scope: scope, - submitJobSuccessCounter: labeled.NewCounter("submit_job_success", "Flink job submission successful", flinkJmClientScope), - submitJobFailureCounter: labeled.NewCounter("submit_job_failure", "Flink job submission failed", flinkJmClientScope), - cancelJobSuccessCounter: labeled.NewCounter("cancel_job_success", "Flink job cancellation successful", flinkJmClientScope), - cancelJobFailureCounter: labeled.NewCounter("cancel_job_failure", "Flink job cancellation failed", flinkJmClientScope), - forceCancelJobSuccessCounter: labeled.NewCounter("force_cancel_job_success", "Flink forced job cancellation successful", flinkJmClientScope), - forceCancelJobFailureCounter: labeled.NewCounter("force_cancel_job_failure", "Flink forced job cancellation failed", flinkJmClientScope), - checkSavepointSuccessCounter: labeled.NewCounter("check_savepoint_status_success", "Flink check savepoint status successful", flinkJmClientScope), - checkSavepointFailureCounter: labeled.NewCounter("check_savepoint_status_failure", "Flink check savepoint status failed", flinkJmClientScope), - getJobsSuccessCounter: labeled.NewCounter("get_jobs_success", "Get flink jobs succeeded", flinkJmClientScope), - getJobsFailureCounter: labeled.NewCounter("get_jobs_failure", "Get flink jobs failed", flinkJmClientScope), - getJobConfigSuccessCounter: labeled.NewCounter("get_job_config_success", "Get flink job config succeeded", flinkJmClientScope), - getJobConfigFailureCounter: labeled.NewCounter("get_job_config_failure", "Get flink job config failed", flinkJmClientScope), - getClusterSuccessCounter: labeled.NewCounter("get_cluster_success", "Get cluster overview succeeded", flinkJmClientScope), - getClusterFailureCounter: labeled.NewCounter("get_cluster_failure", "Get cluster overview failed", flinkJmClientScope), - getCheckpointsSuccessCounter: labeled.NewCounter("get_checkpoints_success", "Get checkpoint request succeeded", flinkJmClientScope), - getCheckpointsFailureCounter: labeled.NewCounter("get_checkpoints_failed", "Get checkpoint request failed", flinkJmClientScope), - savepointJobSuccessCounter: labeled.NewCounter("savepoint_job_success", "Savepoint job request succeeded", flinkJmClientScope), - savepointJobFailureCounter: labeled.NewCounter("savepoint_job_failed", "Savepoint job request failed", flinkJmClientScope), + scope: scope, + submitJobSuccessCounter: labeled.NewCounter("submit_job_success", "Flink job submission successful", flinkJmClientScope), + submitJobFailureCounter: labeled.NewCounter("submit_job_failure", "Flink job submission failed", flinkJmClientScope), + cancelJobSuccessCounter: labeled.NewCounter("cancel_job_success", "Flink job cancellation successful", flinkJmClientScope), + cancelJobFailureCounter: labeled.NewCounter("cancel_job_failure", "Flink job cancellation failed", flinkJmClientScope), + forceCancelJobSuccessCounter: labeled.NewCounter("force_cancel_job_success", "Flink forced job cancellation successful", flinkJmClientScope), + forceCancelJobFailureCounter: labeled.NewCounter("force_cancel_job_failure", "Flink forced job cancellation failed", flinkJmClientScope), + checkSavepointSuccessCounter: labeled.NewCounter("check_savepoint_status_success", "Flink check savepoint status successful", flinkJmClientScope), + checkSavepointFailureCounter: labeled.NewCounter("check_savepoint_status_failure", "Flink check savepoint status failed", flinkJmClientScope), + getJobsSuccessCounter: labeled.NewCounter("get_jobs_success", "Get flink jobs succeeded", flinkJmClientScope), + getJobsFailureCounter: labeled.NewCounter("get_jobs_failure", "Get flink jobs failed", flinkJmClientScope), + getJobConfigSuccessCounter: labeled.NewCounter("get_job_config_success", "Get flink job config succeeded", flinkJmClientScope), + getJobConfigFailureCounter: labeled.NewCounter("get_job_config_failure", "Get flink job config failed", flinkJmClientScope), + getClusterSuccessCounter: labeled.NewCounter("get_cluster_success", "Get cluster overview succeeded", flinkJmClientScope), + getClusterFailureCounter: labeled.NewCounter("get_cluster_failure", "Get cluster overview failed", flinkJmClientScope), + getCheckpointsSuccessCounter: labeled.NewCounter("get_checkpoints_success", "Get checkpoint request succeeded", flinkJmClientScope), + getCheckpointsFailureCounter: labeled.NewCounter("get_checkpoints_failed", "Get checkpoint request failed", flinkJmClientScope), + savepointJobSuccessCounter: labeled.NewCounter("savepoint_job_success", "Savepoint job request succeeded", flinkJmClientScope), + savepointJobFailureCounter: labeled.NewCounter("savepoint_job_failed", "Savepoint job request failed", flinkJmClientScope), + getCheckpointsConfigSuccessCounter: labeled.NewCounter("get_checkpoints_config_success", "Get checkpoint config request succeeded", flinkJmClientScope), + getCheckpointsConfigFailureCounter: labeled.NewCounter("get_checkpoints_config_failed", "Get checkpoint config request failed", flinkJmClientScope), } } @@ -369,6 +374,26 @@ func (c *FlinkJobManagerClient) GetCheckpointCounts(ctx context.Context, url str c.metrics.getCheckpointsSuccessCounter.Inc(ctx) return &checkpointResponse, nil } +func (c *FlinkJobManagerClient) GetCheckpointConfig(ctx context.Context, url string, jobID string) (*CheckpointConfigResponse, error) { + endpoint := fmt.Sprintf(url+checkpointsURL+"/config", jobID) + response, err := c.executeRequest(ctx, httpGet, endpoint, nil) + if err != nil { + c.metrics.getCheckpointsConfigFailureCounter.Inc(ctx) + return nil, GetRetryableError(err, v1beta1.GetCheckpointConfig, GlobalFailure, DefaultRetries) + } + if response != nil && !response.IsSuccess() { + c.metrics.getCheckpointsConfigFailureCounter.Inc(ctx) + return nil, GetRetryableError(err, v1beta1.GetCheckpointConfig, response.Status(), DefaultRetries) + } + + var checkpointConfigResponse CheckpointConfigResponse + if err = json.Unmarshal(response.Body(), &checkpointConfigResponse); err != nil { + logger.Errorf(ctx, "Failed to unmarshal checkpointConfigResponse %v, err %v", response, err) + } + + c.metrics.getCheckpointsSuccessCounter.Inc(ctx) + return &checkpointConfigResponse, nil +} func (c *FlinkJobManagerClient) GetJobOverview(ctx context.Context, url string, jobID string) (*FlinkJobOverview, error) { endpoint := fmt.Sprintf(url+GetJobsOverviewURL, jobID) diff --git a/pkg/controller/flink/client/entities.go b/pkg/controller/flink/client/entities.go index ffbdb8f3..b89caef9 100644 --- a/pkg/controller/flink/client/entities.go +++ b/pkg/controller/flink/client/entities.go @@ -139,12 +139,21 @@ type LatestCheckpoints struct { Restored *CheckpointStatistics `json:"restored,omitempty"` } +type ExternalizedCheckpoints struct { + Enabled bool `json:"enable,omitempty"` + DeleteOnCancellation bool `json:"delete_on_cancellation,omitempty"` +} + type CheckpointResponse struct { Counts map[string]int32 `json:"counts"` Latest LatestCheckpoints `json:"latest"` History []CheckpointStatistics `json:"history"` } +type CheckpointConfigResponse struct { + Externalization *ExternalizedCheckpoints `json:"externalization"` +} + type TaskManagerStats struct { Path string `json:"path"` DataPort int32 `json:"dataPort"` diff --git a/pkg/controller/flink/client/mock/mock_api.go b/pkg/controller/flink/client/mock/mock_api.go index 873d96a8..18db25c1 100644 --- a/pkg/controller/flink/client/mock/mock_api.go +++ b/pkg/controller/flink/client/mock/mock_api.go @@ -16,6 +16,7 @@ type GetLatestCheckpointFunc func(ctx context.Context, url string, jobID string) type GetJobConfigFunc func(ctx context.Context, url string, jobID string) (*client.JobConfigResponse, error) type GetTaskManagersFunc func(ctx context.Context, url string) (*client.TaskManagersResponse, error) type GetCheckpointCountsFunc func(ctx context.Context, url string, jobID string) (*client.CheckpointResponse, error) +type GetCheckpointConfigFunc func(ctx context.Context, url string, jobID string) (*client.CheckpointConfigResponse, error) type GetJobOverviewFunc func(ctx context.Context, url string, jobID string) (*client.FlinkJobOverview, error) type SavepointJobFunc func(ctx context.Context, url string, jobID string) (string, error) type JobManagerClient struct { @@ -29,6 +30,7 @@ type JobManagerClient struct { GetLatestCheckpointFunc GetLatestCheckpointFunc GetTaskManagersFunc GetTaskManagersFunc GetCheckpointCountsFunc GetCheckpointCountsFunc + GetCheckpointConfigFunc GetCheckpointConfigFunc GetJobOverviewFunc GetJobOverviewFunc SavepointJobFunc SavepointJobFunc } @@ -103,6 +105,13 @@ func (m *JobManagerClient) GetCheckpointCounts(ctx context.Context, url string, return nil, nil } +func (m *JobManagerClient) GetCheckpointConfig(ctx context.Context, url string, jobID string) (*client.CheckpointConfigResponse, error) { + if m.GetCheckpointConfigFunc != nil { + return m.GetCheckpointConfigFunc(ctx, url, jobID) + } + return nil, nil +} + func (m *JobManagerClient) GetJobOverview(ctx context.Context, url string, jobID string) (*client.FlinkJobOverview, error) { if m.GetJobOverviewFunc != nil { return m.GetJobOverviewFunc(ctx, url, jobID) diff --git a/pkg/controller/flink/config.go b/pkg/controller/flink/config.go index 1e375b32..c43b7a0c 100644 --- a/pkg/controller/flink/config.go +++ b/pkg/controller/flink/config.go @@ -71,7 +71,6 @@ func getInternalMetricsQueryPort(app *v1beta1.FlinkApplication) int32 { func getMaxCheckpointRestoreAgeSeconds(app *v1beta1.FlinkApplication) int32 { return firstNonNil(app.Spec.MaxCheckpointRestoreAgeSeconds, MaxCheckpointRestoreAgeSeconds) } - func getTaskManagerMemory(application *v1beta1.FlinkApplication) int64 { tmResources := application.Spec.TaskManagerConfig.Resources if tmResources == nil { diff --git a/pkg/controller/flink/flink.go b/pkg/controller/flink/flink.go index 9bb8418f..0e011d7d 100644 --- a/pkg/controller/flink/flink.go +++ b/pkg/controller/flink/flink.go @@ -85,6 +85,9 @@ type ControllerInterface interface { // able to savepoint for some reason. FindExternalizedCheckpoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) + // Ensures that application is configured to externalize and *not* delete checkpoints on cancel. + FindExternalizedCheckpointForSavepoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) + // Logs an event to the FlinkApplication resource and to the operator log LogEvent(ctx context.Context, app *v1beta1.FlinkApplication, eventType string, reason string, message string) @@ -468,8 +471,8 @@ func (f *Controller) DeleteOldResourcesForApp(ctx context.Context, app *v1beta1. return nil } -func (f *Controller) FindExternalizedCheckpoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) { - checkpoint, err := f.flinkClient.GetLatestCheckpoint(ctx, f.getURLFromApp(application, hash), f.GetLatestJobID(ctx, application)) +func (f *Controller) findExternalizedCheckpoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string, checkpointMaxAge int32) (string, error) { + checkpoint, err := f.flinkClient.GetLatestCheckpoint(ctx, f.getURLFromApp(application, hash), application.Status.JobStatus.JobID) var checkpointPath string var checkpointTime int64 if err != nil { @@ -490,12 +493,28 @@ func (f *Controller) FindExternalizedCheckpoint(ctx context.Context, application return "", nil } - if isCheckpointOldToRecover(checkpointTime, getMaxCheckpointRestoreAgeSeconds(application)) { + if isCheckpointOldToRecover(checkpointTime, checkpointMaxAge) { logger.Info(ctx, "Found checkpoint to restore from, but was too old") return "", nil } return checkpointPath, nil + +} + +func (f *Controller) FindExternalizedCheckpoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) { + return f.findExternalizedCheckpoint(ctx, application, hash, getMaxCheckpointRestoreAgeSeconds(application)) +} + +func (f *Controller) FindExternalizedCheckpointForSavepoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) { + checkpointConfig, err := f.flinkClient.GetCheckpointConfig(ctx, f.getURLFromApp(application, hash), application.Status.JobStatus.JobID) + if err != nil { + return "", err + } + if checkpointConfig.Externalization.Enabled && !checkpointConfig.Externalization.DeleteOnCancellation { + return "", fmt.Errorf("Checkpoint configuration not compatable for starting from checkpoints") + } + return f.findExternalizedCheckpoint(ctx, application, hash, getMaxCheckpointRestoreAgeSeconds(application)) } func isCheckpointOldToRecover(checkpointTime int64, maxCheckpointRecoveryAgeSec int32) bool { diff --git a/pkg/controller/flink/flink_test.go b/pkg/controller/flink/flink_test.go index 09b1dfc7..2461f9c3 100644 --- a/pkg/controller/flink/flink_test.go +++ b/pkg/controller/flink/flink_test.go @@ -826,6 +826,16 @@ func TestJobStatusUpdated(t *testing.T) { }, nil } + mockJmClient.GetCheckpointConfigFunc = func(ctx context.Context, url string, jobID string) (*client.CheckpointConfigResponse, error) { + assert.Equal(t, url, "http://app-name-hash.ns:8081") + return &client.CheckpointConfigResponse{ + Externalization: &client.ExternalizedCheckpoints{ + Enabled: true, + DeleteOnCancellation: false, + }, + }, nil + } + flinkApp.Status.JobStatus.JobID = "abc" expectedTime := metaV1.NewTime(time.Unix(startTime/1000, 0)) _, err = flinkControllerForTest.CompareAndUpdateJobStatus(context.Background(), &flinkApp, "hash") diff --git a/pkg/controller/flink/mock/mock_flink.go b/pkg/controller/flink/mock/mock_flink.go index 830bcfa2..fd01e39a 100644 --- a/pkg/controller/flink/mock/mock_flink.go +++ b/pkg/controller/flink/mock/mock_flink.go @@ -22,6 +22,7 @@ type GetJobsForApplicationFunc func(ctx context.Context, application *v1beta1.Fl type GetJobForApplicationFunc func(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (*client.FlinkJobOverview, error) type GetCurrentDeploymentsForAppFunc func(ctx context.Context, application *v1beta1.FlinkApplication) (*common.FlinkDeployment, error) type FindExternalizedCheckpointFunc func(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) +type FindExternalizedCheckpointForSavepointFunc func(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) type CompareAndUpdateClusterStatusFunc func(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (bool, error) type CompareAndUpdateJobStatusFunc func(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (bool, error) type GetLatestClusterStatusFunc func(ctx context.Context, app *v1beta1.FlinkApplication) v1beta1.FlinkClusterStatus @@ -36,34 +37,36 @@ type DeleteStatusPostTeardownFunc func(ctx context.Context, application *v1beta1 type GetJobToDeleteForApplicationFunc func(ctx context.Context, app *v1beta1.FlinkApplication, hash string) (*client.FlinkJobOverview, error) type GetVersionAndJobIDForHashFunc func(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, string, error) type GetVersionAndHashPostTeardownFunc func(ctx context.Context, application *v1beta1.FlinkApplication) (v1beta1.FlinkApplicationVersion, string) + type FlinkController struct { - CreateClusterFunc CreateClusterFunc - DeleteOldResourcesForAppFunc DeleteOldResourcesForApp - SavepointFunc SavepointFunc - ForceCancelFunc ForceCancelFunc - StartFlinkJobFunc StartFlinkJobFunc - GetSavepointStatusFunc GetSavepointStatusFunc - IsClusterReadyFunc IsClusterReadyFunc - IsServiceReadyFunc IsServiceReadyFunc - GetJobsForApplicationFunc GetJobsForApplicationFunc - GetJobForApplicationFunc GetJobForApplicationFunc - GetCurrentDeploymentsForAppFunc GetCurrentDeploymentsForAppFunc - FindExternalizedCheckpointFunc FindExternalizedCheckpointFunc - Events []corev1.Event - CompareAndUpdateClusterStatusFunc CompareAndUpdateClusterStatusFunc - CompareAndUpdateJobStatusFunc CompareAndUpdateJobStatusFunc - GetLatestClusterStatusFunc GetLatestClusterStatusFunc - GetLatestJobStatusFunc GetLatestJobStatusFunc - GetLatestJobIDFunc GetLatestJobIDFunc - UpdateLatestJobIDFunc UpdateLatestJobIDFunc - UpdateLatestJobStatusFunc UpdateLatestJobStatusFunc - UpdateLatestClusterStatusFunc UpdateLatestClusterStatusFunc - UpdateLatestVersionAndHashFunc UpdateLatestVersionAndHashFunc - DeleteResourcesForAppWithHashFunc DeleteResourcesForAppWithHashFunc - DeleteStatusPostTeardownFunc DeleteStatusPostTeardownFunc - GetJobToDeleteForApplicationFunc GetJobToDeleteForApplicationFunc - GetVersionAndJobIDForHashFunc GetVersionAndJobIDForHashFunc - GetVersionAndHashPostTeardownFunc GetVersionAndHashPostTeardownFunc + CreateClusterFunc CreateClusterFunc + DeleteOldResourcesForAppFunc DeleteOldResourcesForApp + SavepointFunc SavepointFunc + ForceCancelFunc ForceCancelFunc + StartFlinkJobFunc StartFlinkJobFunc + GetSavepointStatusFunc GetSavepointStatusFunc + IsClusterReadyFunc IsClusterReadyFunc + IsServiceReadyFunc IsServiceReadyFunc + GetJobsForApplicationFunc GetJobsForApplicationFunc + GetJobForApplicationFunc GetJobForApplicationFunc + GetCurrentDeploymentsForAppFunc GetCurrentDeploymentsForAppFunc + FindExternalizedCheckpointFunc FindExternalizedCheckpointFunc + FindExternalizedCheckpointForSavepointFunc FindExternalizedCheckpointForSavepointFunc + Events []corev1.Event + CompareAndUpdateClusterStatusFunc CompareAndUpdateClusterStatusFunc + CompareAndUpdateJobStatusFunc CompareAndUpdateJobStatusFunc + GetLatestClusterStatusFunc GetLatestClusterStatusFunc + GetLatestJobStatusFunc GetLatestJobStatusFunc + GetLatestJobIDFunc GetLatestJobIDFunc + UpdateLatestJobIDFunc UpdateLatestJobIDFunc + UpdateLatestJobStatusFunc UpdateLatestJobStatusFunc + UpdateLatestClusterStatusFunc UpdateLatestClusterStatusFunc + UpdateLatestVersionAndHashFunc UpdateLatestVersionAndHashFunc + DeleteResourcesForAppWithHashFunc DeleteResourcesForAppWithHashFunc + DeleteStatusPostTeardownFunc DeleteStatusPostTeardownFunc + GetJobToDeleteForApplicationFunc GetJobToDeleteForApplicationFunc + GetVersionAndJobIDForHashFunc GetVersionAndJobIDForHashFunc + GetVersionAndHashPostTeardownFunc GetVersionAndHashPostTeardownFunc } func (m *FlinkController) GetCurrentDeploymentsForApp(ctx context.Context, application *v1beta1.FlinkApplication) (*common.FlinkDeployment, error) { @@ -151,6 +154,13 @@ func (m *FlinkController) FindExternalizedCheckpoint(ctx context.Context, applic return "", nil } +func (m *FlinkController) FindExternalizedCheckpointForSavepoint(ctx context.Context, application *v1beta1.FlinkApplication, hash string) (string, error) { + if m.FindExternalizedCheckpointForSavepointFunc != nil { + return m.FindExternalizedCheckpointFunc(ctx, application, hash) + } + return "", nil +} + func (m *FlinkController) LogEvent(ctx context.Context, app *v1beta1.FlinkApplication, eventType string, reason string, message string) { m.Events = append(m.Events, corev1.Event{ InvolvedObject: corev1.ObjectReference{ diff --git a/pkg/controller/flinkapplication/flink_state_machine.go b/pkg/controller/flinkapplication/flink_state_machine.go index 6426c293..3292cf67 100644 --- a/pkg/controller/flinkapplication/flink_state_machine.go +++ b/pkg/controller/flinkapplication/flink_state_machine.go @@ -233,6 +233,7 @@ func (s *FlinkStateMachine) IsTimeToHandlePhase(application *v1beta1.FlinkApplic // In this state we create a new cluster, either due to an entirely new FlinkApplication or due to an update. func (s *FlinkStateMachine) handleNewOrUpdating(ctx context.Context, application *v1beta1.FlinkApplication) (bool, error) { + // TODO: add up-front validation on the FlinkApplication resource if rollback, reason := s.shouldRollback(ctx, application); rollback { // we've failed to make progress; move to deploy failed @@ -305,9 +306,9 @@ func (s *FlinkStateMachine) handleClusterStarting(ctx context.Context, applicati logger.Infof(ctx, "Flink cluster has started successfully") // TODO: in single mode move to submitting job - if application.Spec.SavepointDisabled && !v1beta1.IsBlueGreenDeploymentMode(application.Status.DeploymentMode) { + if application.Spec.SavepointingDisabled() && !v1beta1.IsBlueGreenDeploymentMode(application.Status.DeploymentMode) { s.updateApplicationPhase(application, v1beta1.FlinkApplicationCancelling) - } else if application.Spec.SavepointDisabled && v1beta1.IsBlueGreenDeploymentMode(application.Status.DeploymentMode) { + } else if application.Spec.SavepointingDisabled() && v1beta1.IsBlueGreenDeploymentMode(application.Status.DeploymentMode) { // Blue Green deployment and no savepoint required implies, we directly transition to submitting job s.updateApplicationPhase(application, v1beta1.FlinkApplicationSubmittingJob) } else { @@ -315,6 +316,24 @@ func (s *FlinkStateMachine) handleClusterStarting(ctx context.Context, applicati } return statusChanged, nil } +func (s *FlinkStateMachine) handleApplicationSavepointingWithCheckpoint(ctx context.Context, application *v1beta1.FlinkApplication) (bool, error) { + checkpointPath, err := s.flinkController.FindExternalizedCheckpointForSavepoint(ctx, application, application.Status.DeployHash) + if err != nil { + return statusUnchanged, err + } + + jobID := s.flinkController.GetLatestJobID(ctx, application) + if err := s.flinkController.ForceCancel(ctx, application, application.Status.DeployHash, jobID); err != nil { + return statusUnchanged, err + } + + s.flinkController.LogEvent(ctx, application, corev1.EventTypeNormal, "CancellingJob", + fmt.Sprintf("Cancelling job job %s with a final checkpoint", jobID)) + application.Status.JobStatus.JobID = "" + application.Status.SavepointPath = checkpointPath + s.updateApplicationPhase(application, v1beta1.FlinkApplicationSubmittingJob) + return statusChanged, nil +} func (s *FlinkStateMachine) initializeAppStatusIfEmpty(application *v1beta1.FlinkApplication) { if v1beta1.IsBlueGreenDeploymentMode(application.Status.DeploymentMode) { @@ -331,6 +350,7 @@ func (s *FlinkStateMachine) initializeAppStatusIfEmpty(application *v1beta1.Flin func (s *FlinkStateMachine) handleApplicationSavepointing(ctx context.Context, application *v1beta1.FlinkApplication) (bool, error) { // we've already savepointed (or this is our first deploy), continue on if application.Status.SavepointPath != "" || application.Status.DeployHash == "" { + logger.Debugf(ctx, "Using SavepointPath: %s", application.Status.SavepointPath) s.updateApplicationPhase(application, v1beta1.FlinkApplicationSubmittingJob) return statusChanged, nil } @@ -342,6 +362,11 @@ func (s *FlinkStateMachine) handleApplicationSavepointing(ctx context.Context, a s.updateApplicationPhase(application, v1beta1.FlinkApplicationRecovering) return statusChanged, nil } + // use of checkpoints in the place of savepoints + if application.Spec.UpdateMode == v1beta1.UpdateModeCheckpoint { + return s.handleApplicationSavepointingWithCheckpoint(ctx, application) + } + cancelFlag := getCancelFlag(application) // we haven't started savepointing yet; do so now // TODO: figure out the idempotence of this