Modifier and Type | Field and Description |
---|---|
protected ExecutionConfig |
Plan.executionConfig
Config object for runtime execution parameters.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
ExecutionConfig.disableClosureCleaner()
Disables the ClosureCleaner.
|
ExecutionConfig |
ExecutionConfig.disableObjectReuse()
Disables reusing objects that Flink internally uses for deserialization and passing
data to user-code functions.
|
ExecutionConfig |
ExecutionConfig.disableSysoutLogging()
Disables the printing of progress update messages to
System.out |
ExecutionConfig |
ExecutionConfig.enableClosureCleaner()
Enables the ClosureCleaner.
|
ExecutionConfig |
ExecutionConfig.enableObjectReuse()
Enables reusing objects that Flink internally uses for deserialization and passing
data to user-code functions.
|
ExecutionConfig |
ExecutionConfig.enableSysoutLogging()
Enables the printing of progress update messages to
System.out |
ExecutionConfig |
Plan.getExecutionConfig()
Gets the execution config object.
|
ExecutionConfig |
ExecutionConfig.setAutoWatermarkInterval(long interval)
Sets the interval of the automatic watermark emission.
|
ExecutionConfig |
ExecutionConfig.setExecutionRetryDelay(long executionRetryDelay)
Deprecated.
This method will be replaced by
setRestartStrategy(org.apache.flink.api.common.restartstrategy.RestartStrategies.RestartStrategyConfiguration) . The
RestartStrategies.FixedDelayRestartStrategyConfiguration contains the delay between
successive execution attempts. |
ExecutionConfig |
ExecutionConfig.setNumberOfExecutionRetries(int numberOfExecutionRetries)
Deprecated.
This method will be replaced by
setRestartStrategy(org.apache.flink.api.common.restartstrategy.RestartStrategies.RestartStrategyConfiguration) . The
RestartStrategies.FixedDelayRestartStrategyConfiguration contains the number of
execution retries. |
ExecutionConfig |
ExecutionConfig.setParallelism(int parallelism)
Sets the parallelism for operations executed through this environment.
|
Modifier and Type | Method and Description |
---|---|
void |
Plan.setExecutionConfig(ExecutionConfig executionConfig)
Sets the runtime config object defining execution parameters.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
RuntimeContext.getExecutionConfig()
Returns the
ExecutionConfig for the currently executing
job. |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
AbstractRuntimeUDFContext.getExecutionConfig() |
Constructor and Description |
---|
AbstractRuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Accumulator<?,?>> accumulators,
Map<String,Future<Path>> cpTasks) |
RuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators) |
Modifier and Type | Method and Description |
---|---|
protected abstract List<OUT> |
SingleInputOperator.executeOnCollections(List<IN> inputData,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected void |
GenericDataSinkBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected abstract List<OUT> |
DualInputOperator.executeOnCollections(List<IN1> inputData1,
List<IN2> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<T> |
Union.executeOnCollections(List<T> inputData1,
List<T> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<OUT> |
GenericDataSourceBase.executeOnCollections(RuntimeContext ctx,
ExecutionConfig executionConfig) |
Constructor and Description |
---|
CollectionExecutor(ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
protected List<IN> |
SortPartitionOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<IN> |
PartitionOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<OUT> |
MapPartitionOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
MapOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
GroupReduceOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
GroupCombineOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
FlatMapOperatorBase.executeOnCollections(List<IN> input,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
OuterJoinOperatorBase.executeOnCollections(List<IN1> leftInput,
List<IN2> rightInput,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<OUT> |
InnerJoinOperatorBase.executeOnCollections(List<IN1> inputData1,
List<IN2> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<OUT> |
CrossOperatorBase.executeOnCollections(List<IN1> inputData1,
List<IN2> inputData2,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
CoGroupRawOperatorBase.executeOnCollections(List<IN1> input1,
List<IN2> input2,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
CoGroupOperatorBase.executeOnCollections(List<IN1> input1,
List<IN2> input2,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<ST> |
DeltaIterationBase.executeOnCollections(List<ST> inputData1,
List<WT> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<T> |
ReduceOperatorBase.executeOnCollections(List<T> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<T> |
FilterOperatorBase.executeOnCollections(List<T> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<T> |
BulkIterationBase.executeOnCollections(List<T> inputData,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
void |
StateDescriptor.initializeSerializerUnlessSet(ExecutionConfig executionConfig)
Initializes the serializer, unless it has been initialized before.
|
Modifier and Type | Method and Description |
---|---|
PrimitiveArrayComparator<T,?> |
PrimitiveArrayTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeComparator<T> |
BasicTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeComparator<T> |
AtomicType.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig)
Creates a comparator for this type.
|
abstract TypeSerializer<T> |
TypeInformation.createSerializer(ExecutionConfig config)
Creates a serializer for the type.
|
TypeSerializer<T> |
PrimitiveArrayTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<Nothing> |
NothingTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
BasicTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
BasicArrayTypeInfo.createSerializer(ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeComparator<T> |
CompositeType.createComparator(int[] logicalKeyFields,
boolean[] orders,
int logicalFieldOffset,
ExecutionConfig config)
Generic implementation of the comparator creation.
|
TypeComparator<T> |
CompositeType.TypeComparatorBuilder.createTypeComparator(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
ExecutionEnvironment.getConfig()
Gets the config object that defines execution parameters.
|
Modifier and Type | Method and Description |
---|---|
void |
TypeSerializerOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
void |
LocalCollectionOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
void |
CsvOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig)
The purpose of this method is solely to check whether the data type to be processed
is in fact a tuple type.
|
Constructor and Description |
---|
PlanProjectOperator(int[] fields,
String name,
TypeInformation<T> inType,
TypeInformation<R> outType,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeComparator<T> |
WritableTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeComparator<T> |
ValueTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeComparator<T> |
GenericTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeComparator<T> |
EnumTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeSerializer<T> |
WritableTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
ValueTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TupleSerializer<T> |
TupleTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
PojoTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<T> |
ObjectArrayTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<InvalidTypesException> |
MissingTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
GenericTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<T> |
EnumTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<Either<L,R>> |
EitherTypeInfo.createSerializer(ExecutionConfig config) |
void |
InputTypeConfigurable.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig)
Method that is called on an
OutputFormat when it is passed to
the DataSet's output method. |
Constructor and Description |
---|
PojoSerializer(Class<T> clazz,
TypeSerializer<?>[] fieldSerializers,
Field[] fields,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
static void |
Serializers.recursivelyRegisterType(Class<?> type,
ExecutionConfig config,
Set<Class<?>> alreadySeen) |
static void |
Serializers.recursivelyRegisterType(TypeInformation<?> typeInfo,
ExecutionConfig config,
Set<Class<?>> alreadySeen) |
Constructor and Description |
---|
KryoSerializer(Class<T> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
ExecutionEnvironment.getConfig()
Gets the config object.
|
Modifier and Type | Method and Description |
---|---|
void |
ScalaCsvOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig)
The purpose of this method is solely to check whether the data type to be processed
is in fact a tuple type.
|
Modifier and Type | Method and Description |
---|---|
TypeComparator<scala.Enumeration.Value> |
EnumValueTypeInfo.createComparator(boolean ascOrder,
ExecutionConfig config) |
TypeSerializer<T> |
TryTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
OptionTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<scala.runtime.Nothing$> |
ScalaNothingTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<scala.Enumeration.Value> |
EnumValueTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
EitherTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<scala.runtime.BoxedUnit> |
UnitTypeInfo.createSerializer(ExecutionConfig config) |
abstract TypeSerializer<T> |
TraversableTypeInfo.createSerializer(ExecutionConfig executionConfig) |
Constructor and Description |
---|
TrySerializer(TypeSerializer<A> elemSerializer,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeComparator<Row> |
RowTypeInfo.createComparator(int[] logicalKeyFields,
boolean[] orders,
int logicalFieldOffset,
ExecutionConfig config) |
TypeComparator<T> |
RenamingProxyTypeInfo.createComparator(int[] logicalKeyFields,
boolean[] orders,
int logicalFieldOffset,
ExecutionConfig executionConfig) |
TypeSerializer<Row> |
RowTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
RenamingProxyTypeInfo.createSerializer(ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
protected List<OUT> |
NoOpBinaryUdfOp.executeOnCollections(List<OUT> inputData1,
List<OUT> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<OUT> |
NoOpUnaryUdfOp.executeOnCollections(List<OUT> inputData,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
static TypeComparatorFactory<?> |
Utils.getShipComparator(Channel channel,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
ExecutionGraph.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
AbstractInvokable.getExecutionConfig()
Returns the global ExecutionConfig, obtained from the job configuration.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
TaskContext.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
static <T> Collector<T> |
BatchTask.initOutputs(AbstractInvokable nepheleTask,
ClassLoader cl,
TaskConfig config,
List<ChainedDriver<?,?>> chainedTasksTarget,
List<RecordWriter<?>> eventualOutputs,
ExecutionConfig executionConfig,
AccumulatorRegistry.Reporter reporter,
Map<String,Accumulator<?,?>> accumulatorMap)
Creates a writer for each output.
|
Modifier and Type | Field and Description |
---|---|
protected ExecutionConfig |
ChainedDriver.executionConfig |
Modifier and Type | Method and Description |
---|---|
void |
ChainedDriver.setup(TaskConfig config,
String taskName,
Collector<OT> outputCollector,
AbstractInvokable parent,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Accumulator<?,?>> accumulatorMap) |
Constructor and Description |
---|
DistributedRuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
DataStream.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
StreamExecutionEnvironment.getConfig()
Gets the config object.
|
Constructor and Description |
---|
ComparableAggregator(int positionToAggregate,
TypeInformation<T> typeInfo,
AggregationFunction.AggregationType aggregationType,
boolean first,
ExecutionConfig config) |
ComparableAggregator(int positionToAggregate,
TypeInformation<T> typeInfo,
AggregationFunction.AggregationType aggregationType,
ExecutionConfig config) |
ComparableAggregator(String field,
TypeInformation<T> typeInfo,
AggregationFunction.AggregationType aggregationType,
boolean first,
ExecutionConfig config) |
SumAggregator(int pos,
TypeInformation<T> typeInfo,
ExecutionConfig config) |
SumAggregator(String field,
TypeInformation<T> typeInfo,
ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
void |
OutputFormatSinkFunction.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
void |
FoldApplyWindowFunction.setOutputType(TypeInformation<ACC> outTypeInfo,
ExecutionConfig executionConfig) |
void |
FoldApplyAllWindowFunction.setOutputType(TypeInformation<ACC> outTypeInfo,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
StreamGraph.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
AbstractStreamOperator.getExecutionConfig()
Gets the execution config defined on the execution environment of the job to which this
operator belongs.
|
Modifier and Type | Method and Description |
---|---|
void |
StreamGroupedFold.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig) |
void |
OutputTypeConfigurable.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig)
Is called by the
org.apache.flink.streaming.api.graph.StreamGraph#addOperator(Integer, StreamOperator, TypeInformation, TypeInformation, String)
method when the StreamGraph is generated. |
void |
AbstractUdfStreamOperator.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
DataStream.executionConfig()
Returns the execution config.
|
ExecutionConfig |
StreamExecutionEnvironment.getConfig()
Gets the config object.
|
ExecutionConfig |
DataStream.getExecutionConfig()
Deprecated.
Use
executionConfig instead. |
Modifier and Type | Method and Description |
---|---|
abstract TypeSerializer<W> |
WindowAssigner.getWindowSerializer(ExecutionConfig executionConfig)
Returns a
TypeSerializer for serializing windows that are assigned by
this WindowAssigner . |
TypeSerializer<TimeWindow> |
TumblingProcessingTimeWindows.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
TumblingEventTimeWindows.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
SlidingProcessingTimeWindows.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
SlidingEventTimeWindows.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<GlobalWindow> |
GlobalWindows.getWindowSerializer(ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
void |
SequenceFileWriter.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
void |
RollingSink.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
void |
WindowOperator.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
void |
NonKeyedWindowOperator.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
static <R,F> FieldAccessor<R,F> |
FieldAccessor.create(int pos,
TypeInformation<R> typeInfo,
ExecutionConfig config) |
static <R,F> FieldAccessor<R,F> |
FieldAccessor.create(String field,
TypeInformation<R> typeInfo,
ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
static <X> KeySelector<X,Tuple> |
KeySelectorUtil.getSelectorForKeys(Keys<X> keys,
TypeInformation<X> typeInfo,
ExecutionConfig executionConfig) |
static <X,K> KeySelector<X,K> |
KeySelectorUtil.getSelectorForOneKey(Keys<X> keys,
Partitioner<K> partitioner,
TypeInformation<X> typeInfo,
ExecutionConfig executionConfig) |
Constructor and Description |
---|
TypeInformationKeyValueSerializationSchema(Class<K> keyClass,
Class<V> valueClass,
ExecutionConfig config)
Creates a new de-/serialization schema for the given types.
|
TypeInformationKeyValueSerializationSchema(TypeInformation<K> keyTypeInfo,
TypeInformation<V> valueTypeInfo,
ExecutionConfig ec)
Creates a new de-/serialization schema for the given types.
|
TypeInformationSerializationSchema(TypeInformation<T> typeInfo,
ExecutionConfig ec)
Creates a new de-/serialization schema for the given type.
|
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.