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.
|
ExecutionConfig |
ExecutionConfig.setTaskCancellationInterval(long interval)
Sets the configuration parameter specifying the interval (in milliseconds)
between consecutive attempts to cancel a running task.
|
ExecutionConfig |
ExecutionConfig.setTaskCancellationTimeout(long timeout)
Sets the timeout (in milliseconds) after which an ongoing task cancellation
is considered failed, leading to a fatal TaskManager error.
|
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,
MetricGroup metrics) |
RuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
MetricGroup metrics) |
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 |
---|---|
TypeComparator<T> |
SqlTimeTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
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> |
SqlTimeTypeInfo.createSerializer(ExecutionConfig executionConfig) |
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<T> |
OptionTypeInfo.createComparator(boolean ascending,
ExecutionConfig executionConfig) |
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) |
TypeSerializer<Row> |
RowTypeInfo.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 |
Environment.getExecutionConfig()
Returns the job specific
ExecutionConfig . |
Modifier and Type | Method and Description |
---|---|
SerializedValue<ExecutionConfig> |
JobInformation.getSerializedExecutionConfig() |
Constructor and Description |
---|
ExecutionGraph(Executor futureExecutor,
Executor ioExecutor,
JobID jobId,
String jobName,
Configuration jobConfig,
SerializedValue<ExecutionConfig> serializedConfig,
scala.concurrent.duration.FiniteDuration timeout,
RestartStrategy restartStrategy,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
Scheduler scheduler,
ClassLoader userClassLoader,
MetricGroup metricGroup) |
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<BlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
Constructor and Description |
---|
ExecutionConfigSummary(ExecutionConfig ec) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<ExecutionConfig> |
JobGraph.getSerializedExecutionConfig()
Returns the
ExecutionConfig |
Modifier and Type | Method and Description |
---|---|
void |
JobGraph.setExecutionConfig(ExecutionConfig executionConfig)
Sets a serialized copy of the passed ExecutionConfig.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
AbstractInvokable.getExecutionConfig()
Returns the global ExecutionConfig.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
TaskContext.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
static <T> Collector<T> |
BatchTask.initOutputs(AbstractInvokable containingTask,
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,
MetricGroup metrics) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
RuntimeEnvironment.getExecutionConfig() |
Constructor and Description |
---|
RuntimeEnvironment(JobID jobId,
JobVertexID jobVertexId,
ExecutionAttemptID executionId,
ExecutionConfig executionConfig,
TaskInfo taskInfo,
Configuration jobConfiguration,
Configuration taskConfiguration,
ClassLoader userCodeClassLoader,
MemoryManager memManager,
IOManager ioManager,
BroadcastVariableManager bcVarManager,
AccumulatorRegistry accumulatorRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
ActorGateway jobManager,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask) |
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 |
ContinuousFileReaderOperator.setOutputType(TypeInformation<OUT> outTypeInfo,
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
StreamGraph.addOperator(Integer, String, 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 |
---|---|
void |
AvroKeyValueSinkWriter.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
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) |
Modifier and Type | Method and Description |
---|---|
void |
InternalSingleValueWindowFunction.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig) |
void |
InternalSingleValueAllWindowFunction.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig) |
void |
InternalIterableWindowFunction.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig) |
void |
InternalIterableAllWindowFunction.setOutputType(TypeInformation<OUT> outTypeInfo,
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.