Modifier and Type | Field and Description |
---|---|
static InflightDataRescalingDescriptor |
InflightDataRescalingDescriptor.NO_RESCALE |
Modifier and Type | Method and Description |
---|---|
InflightDataRescalingDescriptor |
TaskStateSnapshot.getInputRescalingDescriptor()
Returns the input channel mapping for rescaling with in-flight data or
NO_RESCALE . |
InflightDataRescalingDescriptor |
OperatorSubtaskState.getInputRescalingDescriptor() |
InflightDataRescalingDescriptor |
TaskStateSnapshot.getOutputRescalingDescriptor()
Returns the output channel mapping for rescaling with in-flight data or
NO_RESCALE . |
InflightDataRescalingDescriptor |
OperatorSubtaskState.getOutputRescalingDescriptor() |
Modifier and Type | Method and Description |
---|---|
OperatorSubtaskState.Builder |
OperatorSubtaskState.Builder.setInputRescalingDescriptor(InflightDataRescalingDescriptor inputRescalingDescriptor) |
OperatorSubtaskState.Builder |
OperatorSubtaskState.Builder.setOutputRescalingDescriptor(InflightDataRescalingDescriptor outputRescalingDescriptor) |
Modifier and Type | Method and Description |
---|---|
InflightDataRescalingDescriptor |
TaskStateManager.getInputRescalingDescriptor() |
InflightDataRescalingDescriptor |
TaskStateManagerImpl.getInputRescalingDescriptor() |
InflightDataRescalingDescriptor |
TaskStateManager.getOutputRescalingDescriptor() |
InflightDataRescalingDescriptor |
TaskStateManagerImpl.getOutputRescalingDescriptor() |
Modifier and Type | Method and Description |
---|---|
static <T> StreamTaskInput<T> |
StreamTaskNetworkInputFactory.create(CheckpointedInputGate checkpointedInputGate,
TypeSerializer<T> inputSerializer,
IOManager ioManager,
StatusWatermarkValve statusWatermarkValve,
int inputIndex,
InflightDataRescalingDescriptor rescalingDescriptorinflightDataRescalingDescriptor,
java.util.function.Function<Integer,StreamPartitioner<?>> gatePartitioners,
TaskInfo taskInfo)
Factory method for
StreamTaskNetworkInput or RescalingStreamTaskNetworkInput
depending on InflightDataRescalingDescriptor . |
static <IN1,IN2> StreamMultipleInputProcessor |
StreamTwoInputProcessorFactory.create(TaskInvokable ownerTask,
CheckpointedInputGate[] checkpointedInputGates,
IOManager ioManager,
MemoryManager memoryManager,
TaskIOMetricGroup taskIOMetricGroup,
TwoInputStreamOperator<IN1,IN2,?> streamOperator,
WatermarkGauge input1WatermarkGauge,
WatermarkGauge input2WatermarkGauge,
OperatorChain<?,?> operatorChain,
StreamConfig streamConfig,
Configuration taskManagerConfig,
Configuration jobConfig,
ExecutionConfig executionConfig,
ClassLoader userClassloader,
Counter numRecordsIn,
InflightDataRescalingDescriptor inflightDataRescalingDescriptor,
java.util.function.Function<Integer,StreamPartitioner<?>> gatePartitioners,
TaskInfo taskInfo) |
static StreamMultipleInputProcessor |
StreamMultipleInputProcessorFactory.create(TaskInvokable ownerTask,
CheckpointedInputGate[] checkpointedInputGates,
StreamConfig.InputConfig[] configuredInputs,
IOManager ioManager,
MemoryManager memoryManager,
TaskIOMetricGroup ioMetricGroup,
Counter mainOperatorRecordsIn,
MultipleInputStreamOperator<?> mainOperator,
WatermarkGauge[] inputWatermarkGauges,
StreamConfig streamConfig,
Configuration taskManagerConfig,
Configuration jobConfig,
ExecutionConfig executionConfig,
ClassLoader userClassloader,
OperatorChain<?,?> operatorChain,
InflightDataRescalingDescriptor inflightDataRescalingDescriptor,
java.util.function.Function<Integer,StreamPartitioner<?>> gatePartitioners,
TaskInfo taskInfo) |
Constructor and Description |
---|
RescalingStreamTaskNetworkInput(CheckpointedInputGate checkpointedInputGate,
TypeSerializer<T> inputSerializer,
IOManager ioManager,
StatusWatermarkValve statusWatermarkValve,
int inputIndex,
InflightDataRescalingDescriptor inflightDataRescalingDescriptor,
java.util.function.Function<Integer,StreamPartitioner<?>> gatePartitioners,
TaskInfo taskInfo) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.