Class JobStatusObserver<R extends AbstractFlinkResource<?,?>>
- java.lang.Object
-
- org.apache.flink.kubernetes.operator.observer.JobStatusObserver<R>
-
public abstract class JobStatusObserver<R extends AbstractFlinkResource<?,?>> extends java.lang.Object
An observer to observe the job status.
-
-
Field Summary
Fields Modifier and Type Field Description protected FlinkConfigManager
configManager
protected EventRecorder
eventRecorder
static java.lang.String
MISSING_SESSION_JOB_ERR
-
Constructor Summary
Constructors Constructor Description JobStatusObserver(FlinkConfigManager configManager, EventRecorder eventRecorder)
-
Method Summary
All Methods Instance Methods Abstract Methods Concrete Methods Modifier and Type Method Description protected abstract java.util.Optional<org.apache.flink.runtime.client.JobStatusMessage>
filterTargetJob(JobStatus status, java.util.List<org.apache.flink.runtime.client.JobStatusMessage> clusterJobStatuses)
Filter the target job status message by the job list from the cluster.boolean
observe(FlinkResourceContext<R> ctx)
Observe the status of the flink job.protected void
onNoJobsFound(R resource, org.apache.flink.configuration.Configuration config)
Callback when no jobs were found on the cluster.protected abstract void
onTargetJobNotFound(R resource, org.apache.flink.configuration.Configuration config)
Callback when no matching target job was found on a cluster where jobs were found.protected abstract void
onTimeout(FlinkResourceContext<R> ctx)
Callback when list jobs timeout.
-
-
-
Field Detail
-
MISSING_SESSION_JOB_ERR
public static final java.lang.String MISSING_SESSION_JOB_ERR
- See Also:
- Constant Field Values
-
eventRecorder
protected final EventRecorder eventRecorder
-
configManager
protected final FlinkConfigManager configManager
-
-
Constructor Detail
-
JobStatusObserver
public JobStatusObserver(FlinkConfigManager configManager, EventRecorder eventRecorder)
-
-
Method Detail
-
observe
public boolean observe(FlinkResourceContext<R> ctx)
Observe the status of the flink job.- Parameters:
ctx
- Resource context.- Returns:
- If job found return true, otherwise return false.
-
onTargetJobNotFound
protected abstract void onTargetJobNotFound(R resource, org.apache.flink.configuration.Configuration config)
Callback when no matching target job was found on a cluster where jobs were found.- Parameters:
resource
- The Flink resource.config
- Deployed/observe configuration.
-
onNoJobsFound
protected void onNoJobsFound(R resource, org.apache.flink.configuration.Configuration config)
Callback when no jobs were found on the cluster.- Parameters:
resource
- The Flink resource.config
- Deployed/observe configuration.
-
onTimeout
protected abstract void onTimeout(FlinkResourceContext<R> ctx)
Callback when list jobs timeout.- Parameters:
ctx
- Resource context.
-
filterTargetJob
protected abstract java.util.Optional<org.apache.flink.runtime.client.JobStatusMessage> filterTargetJob(JobStatus status, java.util.List<org.apache.flink.runtime.client.JobStatusMessage> clusterJobStatuses)
Filter the target job status message by the job list from the cluster.- Parameters:
status
- the target job status.clusterJobStatuses
- the candidate cluster jobs.- Returns:
- The target job status message. If no matched job found,
Optional.empty()
will be returned.
-
-