Class StreamOneInputProcessor<IN>
- java.lang.Object
-
- org.apache.flink.streaming.runtime.io.StreamOneInputProcessor<IN>
-
- Type Parameters:
IN
- The type of the record that can be read with this record reader.
- All Implemented Interfaces:
Closeable
,AutoCloseable
,AvailabilityProvider
,StreamInputProcessor
@Internal public final class StreamOneInputProcessor<IN> extends Object implements StreamInputProcessor
Input reader forOneInputStreamTask
.
-
-
Nested Class Summary
-
Nested classes/interfaces inherited from interface org.apache.flink.runtime.io.AvailabilityProvider
AvailabilityProvider.AvailabilityHelper
-
-
Field Summary
-
Fields inherited from interface org.apache.flink.runtime.io.AvailabilityProvider
AVAILABLE
-
-
Constructor Summary
Constructors Constructor Description StreamOneInputProcessor(StreamTaskInput<IN> input, PushingAsyncDataInput.DataOutput<IN> output, BoundedMultiInput endOfInputAware)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description void
close()
CompletableFuture<?>
getAvailableFuture()
CompletableFuture<Void>
prepareSnapshot(ChannelStateWriter channelStateWriter, long checkpointId)
DataInputStatus
processInput()
In case of two and more input processors this method must callInputSelectable.nextSelection()
to choose which input to consume from next.-
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
-
Methods inherited from interface org.apache.flink.runtime.io.AvailabilityProvider
isApproximatelyAvailable, isAvailable
-
-
-
-
Constructor Detail
-
StreamOneInputProcessor
public StreamOneInputProcessor(StreamTaskInput<IN> input, PushingAsyncDataInput.DataOutput<IN> output, BoundedMultiInput endOfInputAware)
-
-
Method Detail
-
getAvailableFuture
public CompletableFuture<?> getAvailableFuture()
- Specified by:
getAvailableFuture
in interfaceAvailabilityProvider
- Returns:
- a future that is completed if the respective provider is available.
-
processInput
public DataInputStatus processInput() throws Exception
Description copied from interface:StreamInputProcessor
In case of two and more input processors this method must callInputSelectable.nextSelection()
to choose which input to consume from next.- Specified by:
processInput
in interfaceStreamInputProcessor
- Returns:
- input status to estimate whether more records can be processed immediately or not. If
there are no more records available at the moment and the caller should check finished
state and/or
AvailabilityProvider.getAvailableFuture()
. - Throws:
Exception
-
prepareSnapshot
public CompletableFuture<Void> prepareSnapshot(ChannelStateWriter channelStateWriter, long checkpointId) throws CheckpointException
- Specified by:
prepareSnapshot
in interfaceStreamInputProcessor
- Throws:
CheckpointException
-
close
public void close() throws IOException
- Specified by:
close
in interfaceAutoCloseable
- Specified by:
close
in interfaceCloseable
- Throws:
IOException
-
-