Class StandaloneFlinkService
- java.lang.Object
-
- org.apache.flink.kubernetes.operator.service.AbstractFlinkService
-
- org.apache.flink.kubernetes.operator.service.StandaloneFlinkService
-
- All Implemented Interfaces:
FlinkService
public class StandaloneFlinkService extends AbstractFlinkService
Implementation ofFlinkService
submitting and interacting with Standalone Kubernetes Flink clusters and jobs.
-
-
Nested Class Summary
-
Nested classes/interfaces inherited from interface org.apache.flink.kubernetes.operator.service.FlinkService
FlinkService.CancelResult
-
-
Field Summary
-
Fields inherited from class org.apache.flink.kubernetes.operator.service.AbstractFlinkService
artifactManager, executorService, FIELD_NAME_TOTAL_CPU, FIELD_NAME_TOTAL_MEMORY, kubernetesClient, operatorConfig
-
-
Constructor Summary
Constructors Constructor Description StandaloneFlinkService(io.fabric8.kubernetes.client.KubernetesClient kubernetesClient, ArtifactManager artifactManager, java.util.concurrent.ExecutorService executorService, FlinkOperatorConfiguration operatorConfig)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description FlinkService.CancelResult
cancelJob(FlinkDeployment deployment, SuspendMode suspendMode, org.apache.flink.configuration.Configuration conf)
protected FlinkStandaloneKubeClient
createNamespacedKubeClient(org.apache.flink.configuration.Configuration configuration)
protected void
deleteClusterInternal(java.lang.String namespace, java.lang.String clusterId, org.apache.flink.configuration.Configuration conf, io.fabric8.kubernetes.api.model.DeletionPropagation deletionPropagation)
Delete Flink kubernetes cluster by deleting the kubernetes resources directly.protected void
deployApplicationCluster(JobSpec jobSpec, org.apache.flink.configuration.Configuration conf)
void
deploySessionCluster(org.apache.flink.configuration.Configuration conf)
protected io.fabric8.kubernetes.api.model.PodList
getJmPodList(java.lang.String namespace, java.lang.String clusterId)
boolean
scale(FlinkResourceContext<?> ctx, org.apache.flink.configuration.Configuration deployConfig)
protected void
submitClusterInternal(org.apache.flink.configuration.Configuration conf, Mode mode)
-
Methods inherited from class org.apache.flink.kubernetes.operator.service.AbstractFlinkService
atLeastOneCheckpoint, cancelJob, cancelJobOrError, cancelSessionJob, deleteBlocking, deleteClusterDeployment, deleteDeploymentBlocking, deleteHAData, disposeSavepoint, fetchCheckpointInfo, fetchCheckpointStats, fetchSavepointInfo, getCheckpointInfo, getClusterClient, getClusterInfo, getJmPodList, getJobStatus, getKubernetesClient, getLastCheckpoint, getMetrics, getRestClient, getSocketAddress, isHaMetadataAvailable, isJobManagerPortReady, isJobMissing, isJobTerminated, removeOperatorConfigs, requestJobResult, runJar, savepointJobOrError, submitApplicationCluster, submitJobToSessionCluster, submitSessionCluster, triggerCheckpoint, triggerSavepoint, updateStatusAfterClusterDeletion, uploadJar
-
-
-
-
Constructor Detail
-
StandaloneFlinkService
public StandaloneFlinkService(io.fabric8.kubernetes.client.KubernetesClient kubernetesClient, ArtifactManager artifactManager, java.util.concurrent.ExecutorService executorService, FlinkOperatorConfiguration operatorConfig)
-
-
Method Detail
-
deployApplicationCluster
protected void deployApplicationCluster(JobSpec jobSpec, org.apache.flink.configuration.Configuration conf) throws java.lang.Exception
- Specified by:
deployApplicationCluster
in classAbstractFlinkService
- Throws:
java.lang.Exception
-
deploySessionCluster
public void deploySessionCluster(org.apache.flink.configuration.Configuration conf) throws java.lang.Exception
- Specified by:
deploySessionCluster
in classAbstractFlinkService
- Throws:
java.lang.Exception
-
cancelJob
public FlinkService.CancelResult cancelJob(FlinkDeployment deployment, SuspendMode suspendMode, org.apache.flink.configuration.Configuration conf) throws java.lang.Exception
- Throws:
java.lang.Exception
-
getJmPodList
protected io.fabric8.kubernetes.api.model.PodList getJmPodList(java.lang.String namespace, java.lang.String clusterId)
- Specified by:
getJmPodList
in classAbstractFlinkService
-
createNamespacedKubeClient
@VisibleForTesting protected FlinkStandaloneKubeClient createNamespacedKubeClient(org.apache.flink.configuration.Configuration configuration)
-
submitClusterInternal
protected void submitClusterInternal(org.apache.flink.configuration.Configuration conf, Mode mode) throws org.apache.flink.client.deployment.ClusterDeploymentException
- Throws:
org.apache.flink.client.deployment.ClusterDeploymentException
-
deleteClusterInternal
protected void deleteClusterInternal(java.lang.String namespace, java.lang.String clusterId, org.apache.flink.configuration.Configuration conf, io.fabric8.kubernetes.api.model.DeletionPropagation deletionPropagation)
Description copied from class:AbstractFlinkService
Delete Flink kubernetes cluster by deleting the kubernetes resources directly.- Specified by:
deleteClusterInternal
in classAbstractFlinkService
- Parameters:
namespace
- NamespaceclusterId
- ClusterIdconf
- Configuration of the Flink applicationdeletionPropagation
- Resource deletion propagation policy
-
scale
public boolean scale(FlinkResourceContext<?> ctx, org.apache.flink.configuration.Configuration deployConfig)
-
-