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.enableClosureCleaner()
Enables the ClosureCleaner.
|
ExecutionConfig |
ExecutionConfig.enableObjectReuse()
Enables reusing objects that Flink internally uses for deserialization and passing data to
user-code functions.
|
ExecutionConfig |
Plan.getExecutionConfig()
Gets the execution config object.
|
ExecutionConfig |
ExecutionConfig.setAutoWatermarkInterval(long interval)
Sets the interval of the automatic watermark emission.
|
ExecutionConfig |
ExecutionConfig.setClosureCleanerLevel(ExecutionConfig.ClosureCleanerLevel level)
Configures the closure cleaner.
|
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.setLatencyTrackingInterval(long interval)
Interval for sending latency tracking marks from the sources to the sinks.
|
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.
|
Constructor and Description |
---|
ArchivedExecutionConfig(ExecutionConfig ec) |
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,
UserCodeClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Accumulator<?,?>> accumulators,
Map<String,Future<Path>> cpTasks,
OperatorMetricGroup metrics) |
RuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
OperatorMetricGroup metrics) |
RuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
OperatorMetricGroup metrics,
JobID jobID) |
Modifier and Type | Method and Description |
---|---|
protected void |
GenericDataSinkBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected abstract List<OUT> |
SingleInputOperator.executeOnCollections(List<IN> inputData,
RuntimeContext runtimeContext,
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<OUT> |
FlatMapOperatorBase.executeOnCollections(List<IN> input,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
MapPartitionOperatorBase.executeOnCollections(List<IN> inputData,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
GroupCombineOperatorBase.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<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> |
InnerJoinOperatorBase.executeOnCollections(List<IN1> inputData1,
List<IN2> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<OUT> |
CoGroupOperatorBase.executeOnCollections(List<IN1> input1,
List<IN2> input2,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
CoGroupRawOperatorBase.executeOnCollections(List<IN1> input1,
List<IN2> input2,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
CrossOperatorBase.executeOnCollections(List<IN1> inputData1,
List<IN2> inputData2,
RuntimeContext ctx,
ExecutionConfig executionConfig) |
protected List<OUT> |
OuterJoinOperatorBase.executeOnCollections(List<IN1> leftInput,
List<IN2> rightInput,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<ST> |
DeltaIterationBase.executeOnCollections(List<ST> inputData1,
List<WT> inputData2,
RuntimeContext runtimeContext,
ExecutionConfig executionConfig) |
protected List<T> |
BulkIterationBase.executeOnCollections(List<T> inputData,
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) |
Constructor and Description |
---|
TypeInformationSerializationSchema(TypeInformation<T> typeInfo,
ExecutionConfig ec)
Creates a new de-/serialization schema for the given type.
|
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) |
TypeComparator<T> |
AtomicType.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig)
Creates a comparator for this type.
|
TypeComparator<T> |
BasicTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
PrimitiveArrayComparator<T,?> |
PrimitiveArrayTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeComparator<T> |
LocalTimeTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeSerializer<T> |
SqlTimeTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<Nothing> |
NothingTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
BasicTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
PrimitiveArrayTypeInfo.createSerializer(ExecutionConfig executionConfig) |
abstract TypeSerializer<T> |
TypeInformation.createSerializer(ExecutionConfig config)
Creates a serializer for the type.
|
TypeSerializer<T> |
LocalTimeTypeInfo.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 |
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.
|
void |
LocalCollectionOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
void |
TypeSerializerOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
Constructor and Description |
---|
PlanProjectOperator(int[] fields,
String name,
TypeInformation<T> inType,
TypeInformation<R> outType,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
abstract void |
AvroUtils.addAvroSerializersIfRequired(ExecutionConfig reg,
Class<?> type)
Loads the utility class from
flink-avro and adds Avro-specific serializers. |
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) |
TypeComparator<Row> |
RowTypeInfo.createComparator(int[] logicalKeyFields,
boolean[] orders,
int logicalFieldOffset,
ExecutionConfig config) |
TypeSerializer<Row> |
RowTypeInfo.createLegacySerializer(ExecutionConfig config)
Deprecated.
|
PojoSerializer<T> |
PojoTypeInfo.createPojoSerializer(ExecutionConfig config) |
TypeSerializer<T> |
WritableTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<List<T>> |
ListTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<InvalidTypesException> |
MissingTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
ValueTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
GenericTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<T> |
EnumTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<Either<L,R>> |
EitherTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<T> |
ObjectArrayTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<T> |
PojoTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<Row> |
RowTypeInfo.createSerializer(ExecutionConfig config) |
TupleSerializer<T> |
TupleTypeInfo.createSerializer(ExecutionConfig executionConfig) |
TypeSerializer<Map<K,V>> |
MapTypeInfo.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)
Constructor to create a new
PojoSerializer . |
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) |
Constructor and Description |
---|
PlanGenerator(List<DataSink<?>> sinks,
ExecutionConfig config,
int defaultParallelism,
List<Tuple2<String,DistributedCache.DistributedCacheEntry>> cacheFile,
String jobName) |
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 |
---|---|
TypeSerializer<CompactorRequest> |
CompactorRequestTypeInfo.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
void |
GenericJdbcSinkFunction.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
void |
JdbcOutputFormat.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
void |
JdbcXaSinkFunction.setInputType(TypeInformation<?> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<T> |
PulsarSchemaTypeInformation.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
static <T> PulsarDeserializationSchema<T> |
PulsarDeserializationSchema.flinkTypeInfo(TypeInformation<T> information,
ExecutionConfig config)
Create a PulsarDeserializationSchema by using the given
TypeInformation . |
Constructor and Description |
---|
PulsarTypeInformationWrapper(TypeInformation<T> information,
ExecutionConfig config) |
Constructor and Description |
---|
RocksDBKeyedStateBackend(ClassLoader userCodeClassLoader,
File instanceBasePath,
RocksDBResourceContainer optionsContainer,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
org.rocksdb.RocksDB db,
LinkedHashMap<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
Map<String,HeapPriorityQueueSnapshotRestoreWrapper<?>> registeredPQStates,
int keyGroupPrefixBytes,
CloseableRegistry cancelStreamRegistry,
StreamCompressionDecorator keyGroupCompressionDecorator,
ResourceGuard rocksDBResourceGuard,
RocksDBSnapshotStrategyBase<K,?> checkpointSnapshotStrategy,
RocksDBWriteBatchWrapper writeBatchWrapper,
org.rocksdb.ColumnFamilyHandle defaultColumnFamilyHandle,
RocksDBNativeMetricMonitor nativeMetricMonitor,
SerializedCompositeKeyBuilder<K> sharedRocksKeyBuilder,
PriorityQueueSetFactory priorityQueueFactory,
RocksDbTtlCompactFiltersManager ttlCompactFiltersManager,
InternalKeyContext<K> keyContext,
long writeBatchSize) |
RocksDBKeyedStateBackendBuilder(String operatorIdentifier,
ClassLoader userCodeClassLoader,
File instanceBasePath,
RocksDBResourceContainer optionsContainer,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
ExecutionConfig executionConfig,
LocalRecoveryConfig localRecoveryConfig,
EmbeddedRocksDBStateBackend.PriorityQueueStateType priorityQueueStateType,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
StreamCompressionDecorator keyGroupCompressionDecorator,
CloseableRegistry cancelStreamRegistry) |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<org.apache.avro.generic.GenericRecord> |
GenericRecordAvroTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<T> |
AvroTypeInfo.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
void |
AvroKryoSerializerUtils.addAvroSerializersIfRequired(ExecutionConfig reg,
Class<?> type) |
Modifier and Type | Method and Description |
---|---|
TypeComparator<LongValueWithProperHashCode> |
LongValueWithProperHashCode.LongValueWithProperHashCodeTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeSerializer<LongValueWithProperHashCode> |
LongValueWithProperHashCode.LongValueWithProperHashCodeTypeInfo.createSerializer(ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeComparator<ValueArray<T>> |
ValueArrayTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig) |
TypeSerializer<ValueArray<T>> |
ValueArrayTypeInfo.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 |
QueryableStateClient.getExecutionConfig()
Gets the
ExecutionConfig . |
ExecutionConfig |
QueryableStateClient.setExecutionConfig(ExecutionConfig config)
Replaces the existing
ExecutionConfig (possibly null ), with the provided one. |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<VoidNamespace> |
VoidNamespaceTypeInfo.createSerializer(ExecutionConfig config) |
ExecutionConfig |
QueryableStateClient.setExecutionConfig(ExecutionConfig config)
Replaces the existing
ExecutionConfig (possibly null ), with the provided one. |
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 |
---|
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<PermanentBlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
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 the execution config.
|
JobGraphBuilder |
JobGraphBuilder.setExecutionConfig(ExecutionConfig newExecutionConfig) |
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,
UserCodeClassLoader cl,
TaskConfig config,
List<ChainedDriver<?,?>> chainedTasksTarget,
List<RecordWriter<?>> eventualOutputs,
ExecutionConfig executionConfig,
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,
UserCodeClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Accumulator<?,?>> accumulatorMap) |
Modifier and Type | Method and Description |
---|---|
static <E> ExternalSorterBuilder<E> |
ExternalSorter.newBuilder(MemoryManager memoryManager,
TaskInvokable parentTask,
TypeSerializer<E> serializer,
TypeComparator<E> comparator,
ExecutionConfig executionConfig)
Creates a builder for the
ExternalSorter . |
Constructor and Description |
---|
LargeRecordHandler(TypeSerializer<T> serializer,
TypeComparator<T> comparator,
IOManager ioManager,
MemoryManager memManager,
List<MemorySegment> memory,
TaskInvokable memoryOwner,
int maxFilehandles,
ExecutionConfig executionConfig) |
Constructor and Description |
---|
DistributedRuntimeUDFContext(TaskInfo taskInfo,
UserCodeClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
OperatorMetricGroup metrics,
ExternalResourceInfoProvider externalResourceInfoProvider,
JobID jobID) |
Modifier and Type | Field and Description |
---|---|
protected ExecutionConfig |
AbstractKeyedStateBackendBuilder.executionConfig |
protected ExecutionConfig |
DefaultKeyedStateStore.executionConfig |
protected ExecutionConfig |
DefaultOperatorStateBackendBuilder.executionConfig
The execution configuration.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
DefaultOperatorStateBackend.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<VoidNamespace> |
VoidNamespaceTypeInfo.createSerializer(ExecutionConfig config) |
static StreamCompressionDecorator |
AbstractStateBackend.getCompressionDecorator(ExecutionConfig executionConfig) |
Constructor and Description |
---|
AbstractKeyedStateBackend(TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
CloseableRegistry cancelStreamRegistry,
InternalKeyContext<K> keyContext) |
AbstractKeyedStateBackend(TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
CloseableRegistry cancelStreamRegistry,
StreamCompressionDecorator keyGroupCompressionDecorator,
InternalKeyContext<K> keyContext) |
AbstractKeyedStateBackendBuilder(TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
ClassLoader userCodeClassLoader,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
Collection<KeyedStateHandle> stateHandles,
StreamCompressionDecorator keyGroupCompressionDecorator,
CloseableRegistry cancelStreamRegistry) |
DefaultKeyedStateStore(KeyedStateBackend<?> keyedStateBackend,
ExecutionConfig executionConfig) |
DefaultOperatorStateBackend(ExecutionConfig executionConfig,
CloseableRegistry closeStreamOnCancelRegistry,
Map<String,PartitionableListState<?>> registeredOperatorStates,
Map<String,BackendWritableBroadcastState<?,?>> registeredBroadcastStates,
Map<String,PartitionableListState<?>> accessedStatesByName,
Map<String,BackendWritableBroadcastState<?,?>> accessedBroadcastStatesByName,
SnapshotStrategyRunner<OperatorStateHandle,?> snapshotStrategyRunner) |
DefaultOperatorStateBackendBuilder(ClassLoader userClassloader,
ExecutionConfig executionConfig,
boolean asynchronousSnapshots,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Constructor and Description |
---|
HeapKeyedStateBackend(TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
CloseableRegistry cancelStreamRegistry,
StreamCompressionDecorator keyGroupCompressionDecorator,
Map<String,StateTable<K,?,?>> registeredKVStates,
Map<String,HeapPriorityQueueSnapshotRestoreWrapper<?>> registeredPQStates,
LocalRecoveryConfig localRecoveryConfig,
HeapPriorityQueueSetFactory priorityQueueSetFactory,
org.apache.flink.runtime.state.heap.HeapSnapshotStrategy<K> checkpointStrategy,
SnapshotExecutionType snapshotExecutionType,
org.apache.flink.runtime.state.heap.StateTableFactory<K> stateTableFactory,
InternalKeyContext<K> keyContext) |
HeapKeyedStateBackendBuilder(TaskKvStateRegistry kvStateRegistry,
TypeSerializer<K> keySerializer,
ClassLoader userCodeClassLoader,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
Collection<KeyedStateHandle> stateHandles,
StreamCompressionDecorator keyGroupCompressionDecorator,
LocalRecoveryConfig localRecoveryConfig,
HeapPriorityQueueSetFactory priorityQueueSetFactory,
boolean asynchronousSnapshots,
CloseableRegistry cancelStreamRegistry) |
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,
UserCodeClassLoader userCodeClassLoader,
MemoryManager memManager,
IOManager ioManager,
BroadcastVariableManager bcVarManager,
TaskStateManager taskStateManager,
GlobalAggregateManager aggregateManager,
AccumulatorRegistry accumulatorRegistry,
TaskKvStateRegistry kvStateRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
IndexedInputGate[] inputGates,
TaskEventDispatcher taskEventDispatcher,
CheckpointResponder checkpointResponder,
TaskOperatorEventGateway operatorEventGateway,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask,
ExternalResourceInfoProvider externalResourceInfoProvider) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
StateReaderOperator.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
static <KEY,T,W extends Window,OUT> |
WindowReaderOperator.evictingWindow(WindowReaderFunction<StreamRecord<T>,OUT,KEY,W> readerFunction,
TypeInformation<KEY> keyType,
TypeSerializer<W> windowSerializer,
TypeInformation<T> stateType,
ExecutionConfig config) |
void |
StateReaderOperator.setup(ExecutionConfig executionConfig,
KeyedStateBackend<KEY> keyKeyedStateBackend,
InternalTimeServiceManager<KEY> timerServiceManager,
SavepointRuntimeContext ctx) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
SavepointRuntimeContext.getExecutionConfig() |
ExecutionConfig |
SavepointEnvironment.getExecutionConfig() |
Constructor and Description |
---|
ChangelogKeyedStateBackend(AbstractKeyedStateBackend<K> keyedStateBackend,
String subtaskName,
ExecutionConfig executionConfig,
TtlTimeProvider ttlTimeProvider,
StateChangelogWriter<? extends ChangelogStateHandle> stateChangelogWriter,
Collection<ChangelogStateBackendHandle> initialState,
CheckpointStorageWorkerView checkpointStorageWorkerView) |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<CommittableMessage<CommT>> |
CommittableMessageTypeInfo.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
DataStream.getExecutionConfig() |
Modifier and Type | Field and Description |
---|---|
protected ExecutionConfig |
StreamExecutionEnvironment.config
The execution configuration for this environment.
|
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)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
void |
ContinuousFileReaderOperator.setOutputType(TypeInformation<OUT> outTypeInfo,
ExecutionConfig executionConfig) |
void |
ContinuousFileReaderOperatorFactory.setOutputType(TypeInformation<OUT> type,
ExecutionConfig executionConfig) |
void |
FromElementsFunction.setOutputType(TypeInformation<T> outTypeInfo,
ExecutionConfig executionConfig)
Set element type and re-serialize element if required.
|
Constructor and Description |
---|
ContinuousFileReaderOperatorFactory(InputFormat<OUT,? super T> inputFormat,
TypeInformation<OUT> type,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
StreamGraph.getExecutionConfig() |
Constructor and Description |
---|
StreamGraph(ExecutionConfig executionConfig,
CheckpointConfig checkpointConfig,
SavepointRestoreSettings savepointRestoreSettings) |
StreamGraphGenerator(List<Transformation<?>> transformations,
ExecutionConfig executionConfig,
CheckpointConfig checkpointConfig) |
StreamGraphGenerator(List<Transformation<?>> transformations,
ExecutionConfig executionConfig,
CheckpointConfig checkpointConfig,
ReadableConfig configuration) |
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.
|
ExecutionConfig |
AbstractStreamOperatorV2.getExecutionConfig()
Gets the execution config defined on the execution environment of the job to which this
operator belongs.
|
Constructor and Description |
---|
StreamOperatorStateHandler(StreamOperatorStateContext context,
ExecutionConfig executionConfig,
CloseableRegistry closeableRegistry) |
Modifier and Type | Method and Description |
---|---|
static <K> MultiInputSortingDataInput.SelectableSortingInputs |
MultiInputSortingDataInput.wrapInputs(TaskInvokable containingTask,
StreamTaskInput<Object>[] sortingInputs,
KeySelector<Object,K>[] keySelectors,
TypeSerializer<Object>[] inputSerializers,
TypeSerializer<K> keySerializer,
StreamTaskInput<Object>[] passThroughInputs,
MemoryManager memoryManager,
IOManager ioManager,
boolean objectReuse,
double managedMemoryFraction,
Configuration jobConfiguration,
ExecutionConfig executionConfig) |
Constructor and Description |
---|
SortingDataInput(StreamTaskInput<T> wrappedInput,
TypeSerializer<T> typeSerializer,
TypeSerializer<K> keySerializer,
KeySelector<T,K> keySelector,
MemoryManager memoryManager,
IOManager ioManager,
boolean objectReuse,
double managedMemoryFraction,
Configuration jobConfiguration,
TaskInvokable containingTask,
ExecutionConfig executionConfig) |
Constructor and Description |
---|
BatchExecutionKeyedStateBackend(TypeSerializer<K> keySerializer,
KeyGroupRange keyGroupRange,
ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<byte[]> |
PickledByteArrayTypeInfo.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
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 |
---|
AbstractPerWindowStateStore(KeyedStateBackend<?> keyedStateBackend,
ExecutionConfig executionConfig) |
MergingWindowStateStore(KeyedStateBackend<?> keyedStateBackend,
ExecutionConfig executionConfig) |
PerWindowStateStore(KeyedStateBackend<?> keyedStateBackend,
ExecutionConfig executionConfig) |
WindowOperatorBuilder(WindowAssigner<? super T,W> windowAssigner,
Trigger<? super T,? super W> trigger,
ExecutionConfig config,
TypeInformation<T> inputType,
KeySelector<T,K> keySelector,
TypeInformation<K> keyType) |
Modifier and Type | Method and Description |
---|---|
default ExecutionConfig |
ContainingTaskDetails.getExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<T> |
SingleThreadAccessCheckingTypeInfo.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
static <T> void |
StreamingFunctionUtils.setOutputType(Function userFunction,
TypeInformation<T> outTypeInfo,
ExecutionConfig executionConfig) |
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)
Deprecated.
Creates a new de-/serialization schema for the given type.
|
Modifier and Type | Method and Description |
---|---|
<T,R,F> FieldAccessor<T,F> |
DefaultScalaProductFieldAccessorFactory.createRecursiveProductFieldAccessor(int pos,
TypeInformation<T> typeInfo,
FieldAccessor<R,F> innerAccessor,
ExecutionConfig config) |
<T,R,F> FieldAccessor<T,F> |
ScalaProductFieldAccessorFactory.createRecursiveProductFieldAccessor(int pos,
TypeInformation<T> typeInfo,
FieldAccessor<R,F> innerAccessor,
ExecutionConfig config)
Returns a product
FieldAccessor that does support recursion. |
<T,F> FieldAccessor<T,F> |
DefaultScalaProductFieldAccessorFactory.createSimpleProductFieldAccessor(int pos,
TypeInformation<T> typeInfo,
ExecutionConfig config) |
<T,F> FieldAccessor<T,F> |
ScalaProductFieldAccessorFactory.createSimpleProductFieldAccessor(int pos,
TypeInformation<T> typeInfo,
ExecutionConfig config)
Returns a product
FieldAccessor that does not support recursion. |
static <T,F> FieldAccessor<T,F> |
FieldAccessorFactory.getAccessor(TypeInformation<T> typeInfo,
int pos,
ExecutionConfig config)
Creates a
FieldAccessor for the given field position, which can be used to get and
set the specified field on instances of this type. |
static <T,F> FieldAccessor<T,F> |
FieldAccessorFactory.getAccessor(TypeInformation<T> typeInfo,
String field,
ExecutionConfig config)
Creates a
FieldAccessor for the field that is given by a field expression, which can
be used to get and set the specified field on instances of this type. |
Modifier and Type | Method and Description |
---|---|
CatalogManager.Builder |
CatalogManager.Builder.executionConfig(ExecutionConfig executionConfig) |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<MapView<K,V>> |
MapViewTypeInfo.createSerializer(ExecutionConfig config)
Deprecated.
|
TypeSerializer<ListView<T>> |
ListViewTypeInfo.createSerializer(ExecutionConfig config)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
ExecutionConfig |
DummyStreamExecutionEnvironment.getConfig() |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<CountWindow> |
CountSlidingWindowAssigner.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
SessionWindowAssigner.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<CountWindow> |
CountTumblingWindowAssigner.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
CumulativeWindowAssigner.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
TumblingWindowAssigner.getWindowSerializer(ExecutionConfig executionConfig) |
TypeSerializer<TimeWindow> |
SlidingWindowAssigner.getWindowSerializer(ExecutionConfig executionConfig) |
abstract TypeSerializer<W> |
WindowAssigner.getWindowSerializer(ExecutionConfig executionConfig)
Returns a
TypeSerializer for serializing windows that are assigned by this WindowAssigner . |
Modifier and Type | Method and Description |
---|---|
TypeSerializer<T> |
ExternalTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<SortedMap<K,V>> |
SortedMapTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<StringData> |
StringDataTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<T> |
InternalTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<TimestampData> |
TimestampDataTypeInfo.createSerializer(ExecutionConfig config) |
TypeSerializer<DecimalData> |
DecimalDataTypeInfo.createSerializer(ExecutionConfig config) |
Modifier and Type | Method and Description |
---|---|
RawType<T> |
TypeInformationRawType.resolve(ExecutionConfig config)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
TypeComparator<T> |
TimeIntervalTypeInfo.createComparator(boolean sortOrderAscending,
ExecutionConfig executionConfig)
Deprecated.
|
TypeSerializer<T> |
TimeIntervalTypeInfo.createSerializer(ExecutionConfig config)
Deprecated.
|
TypeSerializer<Timestamp> |
TimeIndicatorTypeInfo.createSerializer(ExecutionConfig executionConfig)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
static <T> InputFormat<T,?> |
PythonTableUtils.getCollectionInputFormat(List<T> data,
TypeInformation<T> dataType,
ExecutionConfig config)
Wrap the unpickled python data with an InputFormat.
|
static InputFormat<Row,?> |
PythonTableUtils.getInputFormat(List<Object[]> data,
TypeInformation<Row> dataType,
ExecutionConfig config)
Wrap the unpickled python data with an InputFormat.
|
Copyright © 2014–2023 The Apache Software Foundation. All rights reserved.