Class SparkStructuredStreamingPipelineResult
java.lang.Object
org.apache.beam.runners.spark.structuredstreaming.SparkStructuredStreamingPipelineResult
- All Implemented Interfaces:
PipelineResult
Result of a pipeline submitted to the
SparkStructuredStreamingRunner. The pipeline runs
asynchronously on a dedicated thread.-
Nested Class Summary
Nested classes/interfaces inherited from interface org.apache.beam.sdk.PipelineResult
PipelineResult.State -
Method Summary
Modifier and TypeMethodDescriptioncancel()Requests cancellation of the pipeline and returns immediately.getState()Retrieves the current state of the pipeline execution.metrics()Returns the object to access metrics from the pipeline.Waits until the pipeline finishes and returns the final status.waitUntilFinish(Duration duration) Waits up todurationfor the execution thread.Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface org.apache.beam.sdk.PipelineResult
drain
-
Method Details
-
getState
Description copied from interface:PipelineResultRetrieves the current state of the pipeline execution.- Specified by:
getStatein interfacePipelineResult- Returns:
- the
PipelineResult.Staterepresenting the state of this pipeline.
-
waitUntilFinish
Description copied from interface:PipelineResultWaits until the pipeline finishes and returns the final status.- Specified by:
waitUntilFinishin interfacePipelineResult- Returns:
- The final state of the pipeline.
-
waitUntilFinish
Waits up todurationfor the execution thread. A pipeline that ends aftercancel()is CANCELLED, any other failure is rethrown and the pipeline is FAILED.- Specified by:
waitUntilFinishin interfacePipelineResult- Parameters:
duration- The time to wait for the pipeline to finish. Provide a value less than 1 ms for an infinite wait.- Returns:
- The final state of the pipeline or null on timeout.
-
metrics
Description copied from interface:PipelineResultReturns the object to access metrics from the pipeline.- Specified by:
metricsin interfacePipelineResult
-
cancel
Requests cancellation of the pipeline and returns immediately.- Specified by:
cancelin interfacePipelineResult- Throws:
IOException- if there is a problem executing the cancel request.
-