Modifier and Type | Method and Description |
---|---|
void |
ChannelStateWriterImpl.addInputData(long checkpointId,
InputChannelInfo info,
int startSeqNum,
CloseableIterator<Buffer> iterator) |
void |
ChannelStateWriter.addInputData(long checkpointId,
InputChannelInfo info,
int startSeqNum,
CloseableIterator<Buffer> data)
Add in-flight buffers from the
InputChannel . |
void |
ChannelStateWriter.NoOpChannelStateWriter.addInputData(long checkpointId,
InputChannelInfo info,
int startSeqNum,
CloseableIterator<Buffer> data) |
Modifier and Type | Method and Description |
---|---|
static void |
NetworkActionsLogger.traceInput(String action,
Buffer buffer,
String taskName,
InputChannelInfo channelInfo,
ChannelStatePersister channelStatePersister,
int sequenceNumber) |
static void |
NetworkActionsLogger.traceRecover(String action,
Buffer buffer,
String taskName,
InputChannelInfo channelInfo) |
Modifier and Type | Field and Description |
---|---|
protected InputChannelInfo |
InputChannel.channelInfo
The info of the input channel to identify it globally within a task.
|
Modifier and Type | Method and Description |
---|---|
InputChannelInfo |
BufferOrEvent.getChannelInfo() |
InputChannelInfo |
InputChannel.getChannelInfo()
Returns the info of this channel, which uniquely identifies the channel in respect to its
operator instance.
|
Modifier and Type | Method and Description |
---|---|
List<InputChannelInfo> |
CheckpointableInput.getChannelInfos() |
List<InputChannelInfo> |
InputGate.getChannelInfos()
Returns the channel infos of this gate.
|
abstract List<InputChannelInfo> |
IndexedInputGate.getUnfinishedChannels()
Returns the list of channels that have not received EndOfPartitionEvent.
|
List<InputChannelInfo> |
SingleInputGate.getUnfinishedChannels() |
Modifier and Type | Method and Description |
---|---|
abstract void |
InputGate.acknowledgeAllRecordsProcessed(InputChannelInfo channelInfo) |
void |
UnionInputGate.acknowledgeAllRecordsProcessed(InputChannelInfo channelInfo) |
void |
SingleInputGate.acknowledgeAllRecordsProcessed(InputChannelInfo channelInfo) |
void |
CheckpointableInput.blockConsumption(InputChannelInfo channelInfo) |
void |
IndexedInputGate.blockConsumption(InputChannelInfo channelInfo) |
void |
CheckpointableInput.resumeConsumption(InputChannelInfo channelInfo) |
abstract void |
InputGate.resumeConsumption(InputChannelInfo channelInfo) |
void |
UnionInputGate.resumeConsumption(InputChannelInfo channelInfo) |
void |
SingleInputGate.resumeConsumption(InputChannelInfo channelInfo) |
void |
BufferOrEvent.setChannelInfo(InputChannelInfo channelInfo) |
Constructor and Description |
---|
BufferOrEvent(AbstractEvent event,
boolean hasPriority,
InputChannelInfo channelInfo,
boolean moreAvailable,
int size,
boolean morePriorityEvents) |
BufferOrEvent(AbstractEvent event,
InputChannelInfo channelInfo) |
BufferOrEvent(Buffer buffer,
InputChannelInfo channelInfo) |
BufferOrEvent(Buffer buffer,
InputChannelInfo channelInfo,
boolean moreAvailable,
boolean morePriorityEvents) |
Constructor and Description |
---|
InputChannelStateHandle(InputChannelInfo info,
StreamStateHandle delegate,
List<Long> offset) |
InputChannelStateHandle(int subtaskIndex,
InputChannelInfo info,
StreamStateHandle delegate,
AbstractChannelStateHandle.StateContentMetaInfo contentMetaInfo) |
InputChannelStateHandle(int subtaskIndex,
InputChannelInfo info,
StreamStateHandle delegate,
List<Long> offset,
long size) |
Modifier and Type | Method and Description |
---|---|
List<InputChannelInfo> |
InputGateWithMetrics.getUnfinishedChannels() |
Modifier and Type | Method and Description |
---|---|
void |
InputGateWithMetrics.acknowledgeAllRecordsProcessed(InputChannelInfo channelInfo) |
void |
InputGateWithMetrics.resumeConsumption(InputChannelInfo channelInfo) |
Modifier and Type | Field and Description |
---|---|
protected Map<InputChannelInfo,Integer> |
AbstractStreamTaskNetworkInput.flattenedChannelIndices |
protected Map<InputChannelInfo,R> |
AbstractStreamTaskNetworkInput.recordDeserializers |
Modifier and Type | Method and Description |
---|---|
List<InputChannelInfo> |
StreamTaskSourceInput.getChannelInfos() |
Modifier and Type | Method and Description |
---|---|
void |
StreamTaskSourceInput.blockConsumption(InputChannelInfo channelInfo) |
protected R |
AbstractStreamTaskNetworkInput.getActiveSerializer(InputChannelInfo channelInfo) |
protected void |
AbstractStreamTaskNetworkInput.releaseDeserializer(InputChannelInfo channelInfo) |
void |
StreamTaskSourceInput.resumeConsumption(InputChannelInfo channelInfo) |
Constructor and Description |
---|
AbstractStreamTaskNetworkInput(CheckpointedInputGate checkpointedInputGate,
TypeSerializer<T> inputSerializer,
StatusWatermarkValve statusWatermarkValve,
int inputIndex,
Map<InputChannelInfo,R> recordDeserializers) |
Modifier and Type | Method and Description |
---|---|
List<InputChannelInfo> |
CheckpointedInputGate.getChannelInfos() |
Modifier and Type | Method and Description |
---|---|
void |
UpstreamRecoveryTracker.handleEndOfRecovery(InputChannelInfo channelInfo) |
protected void |
SingleCheckpointBarrierHandler.markCheckpointAlignedAndTransformState(InputChannelInfo alignedChannel,
CheckpointBarrier barrier,
FunctionWithException<org.apache.flink.streaming.runtime.io.checkpointing.BarrierHandlerState,org.apache.flink.streaming.runtime.io.checkpointing.BarrierHandlerState,Exception> stateTransformer) |
void |
CheckpointBarrierTracker.processBarrier(CheckpointBarrier receivedBarrier,
InputChannelInfo channelInfo,
boolean isRpcTriggered) |
void |
SingleCheckpointBarrierHandler.processBarrier(CheckpointBarrier barrier,
InputChannelInfo channelInfo,
boolean isRpcTriggered) |
abstract void |
CheckpointBarrierHandler.processBarrier(CheckpointBarrier receivedBarrier,
InputChannelInfo channelInfo,
boolean isRpcTriggered) |
void |
CheckpointBarrierTracker.processBarrierAnnouncement(CheckpointBarrier announcedBarrier,
int sequenceNumber,
InputChannelInfo channelInfo) |
void |
SingleCheckpointBarrierHandler.processBarrierAnnouncement(CheckpointBarrier announcedBarrier,
int sequenceNumber,
InputChannelInfo channelInfo) |
abstract void |
CheckpointBarrierHandler.processBarrierAnnouncement(CheckpointBarrier announcedBarrier,
int sequenceNumber,
InputChannelInfo channelInfo) |
void |
CheckpointBarrierTracker.processCancellationBarrier(CancelCheckpointMarker cancelBarrier,
InputChannelInfo channelInfo) |
void |
SingleCheckpointBarrierHandler.processCancellationBarrier(CancelCheckpointMarker cancelBarrier,
InputChannelInfo channelInfo) |
abstract void |
CheckpointBarrierHandler.processCancellationBarrier(CancelCheckpointMarker cancelBarrier,
InputChannelInfo channelInfo) |
void |
CheckpointBarrierTracker.processEndOfPartition(InputChannelInfo channelInfo) |
void |
SingleCheckpointBarrierHandler.processEndOfPartition(InputChannelInfo channelInfo) |
abstract void |
CheckpointBarrierHandler.processEndOfPartition(InputChannelInfo channelInfo) |
Modifier and Type | Method and Description |
---|---|
protected org.apache.flink.streaming.runtime.io.recovery.DemultiplexingRecordDeserializer<T> |
RescalingStreamTaskNetworkInput.getActiveSerializer(InputChannelInfo channelInfo) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.