Class SessionReconciler
- java.lang.Object
-
- org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>
-
- org.apache.flink.kubernetes.operator.reconciler.deployment.SessionReconciler
-
- All Implemented Interfaces:
Reconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment>
public class SessionReconciler extends AbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>
Reconciler responsible for handling the session cluster lifecycle according to the desired and current states.
-
-
Field Summary
Fields Modifier and Type Field Description protected FlinkServiceflinkService-
Fields inherited from class org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractFlinkResourceReconciler
clock, configManager, eventRecorder, kubernetesClient, MSG_ROLLBACK, MSG_SPEC_CHANGED, MSG_SUBMIT, MSG_SUSPENDED, statusRecorder
-
-
Constructor Summary
Constructors Constructor Description SessionReconciler(io.fabric8.kubernetes.client.KubernetesClient kubernetesClient, FlinkService flinkService, FlinkConfigManager configManager, EventRecorder eventRecorder, StatusRecorder<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus> statusRecorder)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description io.javaoperatorsdk.operator.api.reconciler.DeleteControlcleanupInternal(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> context)Shut down and clean up all Flink job/cluster resources.protected voiddeploy(org.apache.flink.kubernetes.operator.api.FlinkDeployment cr, org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec spec, org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus status, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration deployConfig, java.util.Optional<java.lang.String> savepoint, boolean requireHaMetadata)Deploys the target resource spec to Kubernetes.protected org.apache.flink.configuration.ConfigurationgetDeployConfig(io.fabric8.kubernetes.api.model.ObjectMeta meta, org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec spec, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx)Get Flink configuration object for deploying the given spec usingAbstractFlinkResourceReconciler.deploy(CR, SPEC, STATUS, io.javaoperatorsdk.operator.api.reconciler.Context<?>, org.apache.flink.configuration.Configuration, java.util.Optional<java.lang.String>, boolean).protected FlinkServicegetFlinkService(org.apache.flink.kubernetes.operator.api.FlinkDeployment resource, io.javaoperatorsdk.operator.api.reconciler.Context<?> context)Get the Flink service related to the resource and context.protected org.apache.flink.configuration.ConfigurationgetObserveConfig(org.apache.flink.kubernetes.operator.api.FlinkDeployment resource, io.javaoperatorsdk.operator.api.reconciler.Context<?> context)Get Flink configuration for client interactions with the running Flink deployment/session job.protected booleanreadyToReconcile(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration deployConfig)Check whether the given Flink resource is ready to be reconciled or we are still waiting for any pending operation or condition first.booleanreconcileOtherChanges(org.apache.flink.kubernetes.operator.api.FlinkDeployment flinkApp, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration observeConfig)Reconcile any other changes required for this resource that are specific to the reconciler implementation.protected voidreconcileSpecChange(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration observeConfig, org.apache.flink.configuration.Configuration deployConfig, org.apache.flink.kubernetes.operator.api.diff.DiffType type)Reconcile spec upgrade on the currently deployed/suspended Flink resource and update the status accordingly.protected voidrollback(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration observeConfig)Rollback deployed resource to the last stable spec.-
Methods inherited from class org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractFlinkResourceReconciler
cleanup, flinkVersionChanged, reconcile, setClock, setOwnerReference, shouldRecoverDeployment
-
-
-
-
Field Detail
-
flinkService
protected final FlinkService flinkService
-
-
Constructor Detail
-
SessionReconciler
public SessionReconciler(io.fabric8.kubernetes.client.KubernetesClient kubernetesClient, FlinkService flinkService, FlinkConfigManager configManager, EventRecorder eventRecorder, StatusRecorder<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus> statusRecorder)
-
-
Method Detail
-
getFlinkService
protected FlinkService getFlinkService(org.apache.flink.kubernetes.operator.api.FlinkDeployment resource, io.javaoperatorsdk.operator.api.reconciler.Context<?> context)
Description copied from class:AbstractFlinkResourceReconcilerGet the Flink service related to the resource and context.- Specified by:
getFlinkServicein classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
resource- Resource being reconciled.context- Current context.- Returns:
- Flink service implementation.
-
getDeployConfig
protected org.apache.flink.configuration.Configuration getDeployConfig(io.fabric8.kubernetes.api.model.ObjectMeta meta, org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec spec, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx)Description copied from class:AbstractFlinkResourceReconcilerGet Flink configuration object for deploying the given spec usingAbstractFlinkResourceReconciler.deploy(CR, SPEC, STATUS, io.javaoperatorsdk.operator.api.reconciler.Context<?>, org.apache.flink.configuration.Configuration, java.util.Optional<java.lang.String>, boolean).- Specified by:
getDeployConfigin classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
meta- ObjectMeta of the related resource.spec- Spec for which the config should be created.ctx- Reconciliation context.- Returns:
- Deployment configuration.
-
getObserveConfig
protected org.apache.flink.configuration.Configuration getObserveConfig(org.apache.flink.kubernetes.operator.api.FlinkDeployment resource, io.javaoperatorsdk.operator.api.reconciler.Context<?> context)Description copied from class:AbstractFlinkResourceReconcilerGet Flink configuration for client interactions with the running Flink deployment/session job.- Specified by:
getObserveConfigin classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
resource- Related Flink resource.context- Reconciliation context.- Returns:
- Observe configuration.
-
readyToReconcile
protected boolean readyToReconcile(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration deployConfig)Description copied from class:AbstractFlinkResourceReconcilerCheck whether the given Flink resource is ready to be reconciled or we are still waiting for any pending operation or condition first.- Specified by:
readyToReconcilein classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
deployment- Related Flink resource.ctx- Reconciliation context.deployConfig- Deployment configuration.- Returns:
- True if the resource is ready to be reconciled.
-
reconcileSpecChange
protected void reconcileSpecChange(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration observeConfig, org.apache.flink.configuration.Configuration deployConfig, org.apache.flink.kubernetes.operator.api.diff.DiffType type) throws java.lang.ExceptionDescription copied from class:AbstractFlinkResourceReconcilerReconcile spec upgrade on the currently deployed/suspended Flink resource and update the status accordingly.- Specified by:
reconcileSpecChangein classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
deployment- Related Flink resource.observeConfig- Observe configuration.deployConfig- Deployment configuration.- Throws:
java.lang.Exception- Error during spec upgrade.
-
deploy
protected void deploy(org.apache.flink.kubernetes.operator.api.FlinkDeployment cr, org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec spec, org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus status, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration deployConfig, java.util.Optional<java.lang.String> savepoint, boolean requireHaMetadata) throws java.lang.ExceptionDescription copied from class:AbstractFlinkResourceReconcilerDeploys the target resource spec to Kubernetes.- Specified by:
deployin classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
cr- Related resource.spec- Spec that should be deployed to Kubernetes.status- Status object of the resourcectx- Reconciliation context.deployConfig- Flink conf for the deployment.savepoint- Optional savepoint path for applications and session jobs.requireHaMetadata- Flag used by application deployments to validate HA metadata- Throws:
java.lang.Exception- Error during deployment.
-
rollback
protected void rollback(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration observeConfig) throws java.lang.ExceptionDescription copied from class:AbstractFlinkResourceReconcilerRollback deployed resource to the last stable spec.- Specified by:
rollbackin classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
deployment- Related Flink resource.ctx- Reconciliation context.observeConfig- Observe configuration.- Throws:
java.lang.Exception- Error during rollback.
-
reconcileOtherChanges
public boolean reconcileOtherChanges(org.apache.flink.kubernetes.operator.api.FlinkDeployment flinkApp, io.javaoperatorsdk.operator.api.reconciler.Context<?> ctx, org.apache.flink.configuration.Configuration observeConfig) throws java.lang.ExceptionDescription copied from class:AbstractFlinkResourceReconcilerReconcile any other changes required for this resource that are specific to the reconciler implementation.- Specified by:
reconcileOtherChangesin classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
flinkApp- Related Flink resource.observeConfig- Observe configuration.- Returns:
- True if any further reconciliation action was taken.
- Throws:
java.lang.Exception- Error during reconciliation.
-
cleanupInternal
public io.javaoperatorsdk.operator.api.reconciler.DeleteControl cleanupInternal(org.apache.flink.kubernetes.operator.api.FlinkDeployment deployment, io.javaoperatorsdk.operator.api.reconciler.Context<?> context)Description copied from class:AbstractFlinkResourceReconcilerShut down and clean up all Flink job/cluster resources.- Specified by:
cleanupInternalin classAbstractFlinkResourceReconciler<org.apache.flink.kubernetes.operator.api.FlinkDeployment,org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec,org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus>- Parameters:
deployment- Resource being reconciled.context- Current context.- Returns:
- DeleteControl object.
-
-