Modifier and Type | Method and Description |
---|---|
MetricGroup |
TimestampAssignerSupplier.Context.getMetricGroup()
Returns the metric group for the context in which the created
TimestampAssigner
is used. |
MetricGroup |
WatermarkGeneratorSupplier.Context.getMetricGroup()
Returns the metric group for the context in which the created
WatermarkGenerator
is used. |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
RuntimeContext.getMetricGroup()
Returns the metric group for this parallel subtask.
|
Modifier and Type | Method and Description |
---|---|
MetricGroup |
AbstractRuntimeUDFContext.getMetricGroup() |
Constructor and Description |
---|
AbstractRuntimeUDFContext(TaskInfo taskInfo,
UserCodeClassLoader 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 |
---|---|
MetricGroup |
DeserializationSchema.InitializationContext.getMetricGroup()
Returns the metric group for the parallel subtask of the source that runs this
DeserializationSchema . |
MetricGroup |
SerializationSchema.InitializationContext.getMetricGroup()
Returns the metric group for the parallel subtask of the source that runs this
SerializationSchema . |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
Sink.InitContext.metricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
SplitEnumeratorContext.metricGroup() |
MetricGroup |
SourceReaderContext.metricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
TestingSplitEnumeratorContext.metricGroup() |
MetricGroup |
TestingReaderContext.metricGroup() |
Modifier and Type | Method and Description |
---|---|
<K> AbstractKeyedStateBackend<K> |
RocksDBStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
<K> AbstractKeyedStateBackend<K> |
RocksDBStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry,
double managedMemoryFraction) |
Constructor and Description |
---|
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,
RocksDBStateBackend.PriorityQueueStateType priorityQueueStateType,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
StreamCompressionDecorator keyGroupCompressionDecorator,
CloseableRegistry cancelStreamRegistry) |
RocksDBNativeMetricMonitor(RocksDBNativeMetricOptions options,
MetricGroup metricGroup,
org.rocksdb.RocksDB rocksDB) |
Modifier and Type | Field and Description |
---|---|
protected MetricGroup |
AbstractRocksDBRestoreOperation.metricGroup |
Constructor and Description |
---|
AbstractRocksDBRestoreOperation(KeyGroupRange keyGroupRange,
int keyGroupPrefixBytes,
int numberOfTransferringThreads,
CloseableRegistry cancelStreamRegistry,
ClassLoader userCodeClassLoader,
Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
StateSerializerProvider<K> keySerializerProvider,
File instanceBasePath,
File instanceRocksDBPath,
org.rocksdb.DBOptions dbOptions,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
RocksDBNativeMetricOptions nativeMetricOptions,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
RocksDbTtlCompactFiltersManager ttlCompactFiltersManager,
Long writeBufferManagerCapacity) |
RocksDBFullRestoreOperation(KeyGroupRange keyGroupRange,
int keyGroupPrefixBytes,
int numberOfTransferringThreads,
CloseableRegistry cancelStreamRegistry,
ClassLoader userCodeClassLoader,
Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
StateSerializerProvider<K> keySerializerProvider,
File instanceBasePath,
File instanceRocksDBPath,
org.rocksdb.DBOptions dbOptions,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
RocksDBNativeMetricOptions nativeMetricOptions,
MetricGroup metricGroup,
Collection<KeyedStateHandle> restoreStateHandles,
RocksDbTtlCompactFiltersManager ttlCompactFiltersManager,
long writeBatchSize,
Long writeBufferManagerCapacity,
PriorityQueueFlag queueRestoreEnabled) |
RocksDBIncrementalRestoreOperation(String operatorIdentifier,
KeyGroupRange keyGroupRange,
int keyGroupPrefixBytes,
int numberOfTransferringThreads,
CloseableRegistry cancelStreamRegistry,
ClassLoader userCodeClassLoader,
Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
StateSerializerProvider<K> keySerializerProvider,
File instanceBasePath,
File instanceRocksDBPath,
org.rocksdb.DBOptions dbOptions,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
RocksDBNativeMetricOptions nativeMetricOptions,
MetricGroup metricGroup,
Collection<KeyedStateHandle> restoreStateHandles,
RocksDbTtlCompactFiltersManager ttlCompactFiltersManager,
long writeBatchSize,
Long writeBufferManagerCapacity,
PriorityQueueFlag queueRestoreEnabled) |
RocksDBNoneRestoreOperation(KeyGroupRange keyGroupRange,
int keyGroupPrefixBytes,
int numberOfTransferringThreads,
CloseableRegistry cancelStreamRegistry,
ClassLoader userCodeClassLoader,
Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
StateSerializerProvider<K> keySerializerProvider,
File instanceBasePath,
File instanceRocksDBPath,
org.rocksdb.DBOptions dbOptions,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
RocksDBNativeMetricOptions nativeMetricOptions,
MetricGroup metricGroup,
Collection<KeyedStateHandle> restoreStateHandles,
RocksDbTtlCompactFiltersManager ttlCompactFiltersManager,
Long writeBufferManagerCapacity) |
Modifier and Type | Method and Description |
---|---|
void |
ScheduledDropwizardReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
ScheduledDropwizardReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
MetricGroup.addGroup(int name)
Creates a new MetricGroup and adds it to this groups sub-groups.
|
MetricGroup |
MetricGroup.addGroup(String name)
Creates a new MetricGroup and adds it to this groups sub-groups.
|
MetricGroup |
MetricGroup.addGroup(String key,
String value)
Creates a new key-value MetricGroup pair.
|
Modifier and Type | Method and Description |
---|---|
void |
DatadogHttpReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
DatadogHttpReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group) |
Modifier and Type | Class and Description |
---|---|
class |
UnregisteredMetricsGroup
A special
MetricGroup that does not register any metrics at the metrics registry and any
reporters. |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
UnregisteredMetricsGroup.addGroup(int name) |
MetricGroup |
UnregisteredMetricsGroup.addGroup(String name) |
MetricGroup |
UnregisteredMetricsGroup.addGroup(String key,
String value) |
Modifier and Type | Method and Description |
---|---|
void |
JMXReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
JMXReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractPrometheusReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
AbstractPrometheusReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group) |
Modifier and Type | Method and Description |
---|---|
void |
MetricReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group)
Called when a new
Metric was added. |
void |
AbstractReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
MetricReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group)
Called when a
Metric was removed. |
void |
AbstractReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group) |
Constructor and Description |
---|
FlinkMetricContainer(MetricGroup metricGroup) |
Constructor and Description |
---|
CheckpointStatsTracker(int numRememberedCheckpoints,
List<ExecutionJobVertex> jobVertices,
CheckpointCoordinatorConfiguration jobCheckpointingConfiguration,
MetricGroup metricGroup)
Creates a new checkpoint stats tracker.
|
Modifier and Type | Method and Description |
---|---|
static ExecutionGraph |
ExecutionGraphBuilder.buildGraph(ExecutionGraph prior,
JobGraph jobGraph,
Configuration jobManagerConfig,
ScheduledExecutorService futureExecutor,
Executor ioExecutor,
SlotProvider slotProvider,
ClassLoader classLoader,
CheckpointRecoveryFactory recoveryFactory,
Time rpcTimeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
BlobWriter blobWriter,
Time allocationTimeout,
org.slf4j.Logger log,
ShuffleMaster<?> shuffleMaster,
JobMasterPartitionTracker partitionTracker,
FailoverStrategy.Factory failoverStrategyFactory,
ExecutionDeploymentListener executionDeploymentListener,
ExecutionStateUpdateListener executionStateUpdateListener,
long initializationTimestamp) |
static ExecutionGraph |
ExecutionGraphBuilder.buildGraph(ExecutionGraph prior,
JobGraph jobGraph,
Configuration jobManagerConfig,
ScheduledExecutorService futureExecutor,
Executor ioExecutor,
SlotProvider slotProvider,
ClassLoader classLoader,
CheckpointRecoveryFactory recoveryFactory,
Time rpcTimeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
BlobWriter blobWriter,
Time allocationTimeout,
org.slf4j.Logger log,
ShuffleMaster<?> shuffleMaster,
JobMasterPartitionTracker partitionTracker,
long initializationTimestamp)
Builds the ExecutionGraph from the JobGraph.
|
Modifier and Type | Method and Description |
---|---|
void |
FailoverStrategy.registerMetrics(MetricGroup metricGroup)
Tells the FailoverStrategy to register its metrics.
|
Modifier and Type | Method and Description |
---|---|
ShuffleIOOwnerContext |
NettyShuffleEnvironment.createShuffleIOOwnerContext(String ownerName,
ExecutionAttemptID executionAttemptID,
MetricGroup parentGroup) |
void |
NettyShuffleEnvironment.registerLegacyNetworkMetrics(MetricGroup metricGroup,
ResultPartitionWriter[] producedPartitions,
InputGate[] inputGates)
Deprecated.
should be removed in future
|
Modifier and Type | Method and Description |
---|---|
static MetricGroup |
NettyShuffleMetricFactory.createShuffleIOOwnerMetricGroup(MetricGroup parentGroup) |
Modifier and Type | Method and Description |
---|---|
static MetricGroup |
NettyShuffleMetricFactory.createShuffleIOOwnerMetricGroup(MetricGroup parentGroup) |
static void |
NettyShuffleMetricFactory.registerInputMetrics(boolean isDetailedMetrics,
MetricGroup inputGroup,
SingleInputGate[] inputGates) |
static void |
NettyShuffleMetricFactory.registerLegacyNetworkMetrics(boolean isDetailedMetrics,
MetricGroup metricGroup,
ResultPartitionWriter[] producedPartitions,
InputGate[] inputGates)
Deprecated.
should be removed in future
|
static void |
NettyShuffleMetricFactory.registerOutputMetrics(boolean isDetailedMetrics,
MetricGroup outputGroup,
ResultPartition[] resultPartitions) |
static void |
ResultPartitionMetrics.registerQueueLengthMetrics(MetricGroup parent,
ResultPartition[] partitions) |
static void |
InputGateMetrics.registerQueueLengthMetrics(MetricGroup parent,
SingleInputGate[] gates) |
static void |
NettyShuffleMetricFactory.registerShuffleMetrics(MetricGroup metricGroup,
NetworkBufferPool networkBufferPool) |
Constructor and Description |
---|
InputChannelMetrics(MetricGroup... parents) |
Modifier and Type | Method and Description |
---|---|
DistributedRuntimeUDFContext |
AbstractIterativeTask.createRuntimeContext(MetricGroup metrics) |
Modifier and Type | Class and Description |
---|---|
class |
ProxyMetricGroup<P extends MetricGroup>
Metric group which forwards all registration calls to its parent metric group.
|
Modifier and Type | Class and Description |
---|---|
class |
AbstractMetricGroup<A extends AbstractMetricGroup<?>>
Abstract
MetricGroup that contains key functionality for adding metrics and groups. |
class |
ComponentMetricGroup<P extends AbstractMetricGroup<?>>
Abstract
MetricGroup for system components (e.g., TaskManager,
Job, Task, Operator). |
class |
FrontMetricGroup<P extends AbstractMetricGroup<?>>
Metric group which forwards all registration calls to a variable parent metric group that injects
a variable reporter index into calls to
getMetricIdentifier(String) or getMetricIdentifier(String, CharacterFilter) . |
class |
GenericKeyMetricGroup
A
GenericMetricGroup for representing the key part of a key-value metric group pair. |
class |
GenericMetricGroup
A simple named
MetricGroup that is used to hold subgroups of
metrics. |
class |
GenericValueMetricGroup
A
GenericMetricGroup for representing the value part of a key-value metric group pair. |
class |
JobManagerJobMetricGroup
Special
MetricGroup representing everything belonging to a
specific job, running on the JobManager. |
class |
JobManagerMetricGroup
Special
MetricGroup representing a JobManager. |
class |
JobMetricGroup<C extends ComponentMetricGroup<C>>
Special abstract
MetricGroup representing everything belonging
to a specific job. |
class |
OperatorIOMetricGroup
Metric group that contains shareable pre-defined IO-related metrics.
|
class |
OperatorMetricGroup
Special
MetricGroup representing an Operator. |
class |
ProcessMetricGroup
AbstractImitatingJobManagerMetricGroup implementation for process related metrics. |
class |
ProxyMetricGroup<P extends MetricGroup>
Metric group which forwards all registration calls to its parent metric group.
|
class |
ResourceManagerMetricGroup
Metric group which is used by the
ResourceManager to register metrics. |
class |
SlotManagerMetricGroup
Metric group which is used by the
SlotManager to register metrics. |
class |
TaskIOMetricGroup
Metric group that contains shareable pre-defined IO-related metrics.
|
class |
TaskManagerJobMetricGroup
Special
MetricGroup representing everything belonging to a
specific job, running on the TaskManager. |
class |
TaskManagerMetricGroup
Special
MetricGroup representing a TaskManager. |
class |
TaskMetricGroup
Special
MetricGroup representing a Flink runtime Task. |
static class |
UnregisteredMetricGroups.UnregisteredJobManagerJobMetricGroup
A safe drop-in replacement for
JobManagerJobMetricGroup s. |
static class |
UnregisteredMetricGroups.UnregisteredJobManagerMetricGroup
A safe drop-in replacement for
JobManagerMetricGroup s. |
static class |
UnregisteredMetricGroups.UnregisteredOperatorMetricGroup
A safe drop-in replacement for
OperatorMetricGroup s. |
static class |
UnregisteredMetricGroups.UnregisteredProcessMetricGroup
A safe drop-in replacement for
ProcessMetricGroups . |
static class |
UnregisteredMetricGroups.UnregisteredResourceManagerMetricGroup
A safe drop-in replacement for
ResourceManagerMetricGroups . |
static class |
UnregisteredMetricGroups.UnregisteredSlotManagerMetricGroup
A safe drop-in replacement for
SlotManagerMetricGroups . |
static class |
UnregisteredMetricGroups.UnregisteredTaskManagerJobMetricGroup
A safe drop-in replacement for
TaskManagerJobMetricGroup s. |
static class |
UnregisteredMetricGroups.UnregisteredTaskManagerMetricGroup
A safe drop-in replacement for
TaskManagerMetricGroup s. |
static class |
UnregisteredMetricGroups.UnregisteredTaskMetricGroup
A safe drop-in replacement for
TaskMetricGroup s. |
Modifier and Type | Field and Description |
---|---|
protected P |
ProxyMetricGroup.parentMetricGroup |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
ProxyMetricGroup.addGroup(int name) |
MetricGroup |
AbstractMetricGroup.addGroup(int name) |
MetricGroup |
ProxyMetricGroup.addGroup(String name) |
MetricGroup |
AbstractMetricGroup.addGroup(String name) |
MetricGroup |
ProxyMetricGroup.addGroup(String key,
String value) |
MetricGroup |
AbstractMetricGroup.addGroup(String key,
String value) |
MetricGroup |
GenericKeyMetricGroup.addGroup(String key,
String value) |
Modifier and Type | Method and Description |
---|---|
static Tuple2<TaskManagerMetricGroup,MetricGroup> |
MetricUtils.instantiateTaskManagerMetricGroup(MetricRegistry metricRegistry,
String hostName,
ResourceID resourceID,
Optional<Time> systemResourceProbeInterval) |
Modifier and Type | Method and Description |
---|---|
static void |
MetricUtils.instantiateFlinkMemoryMetricGroup(MetricGroup parentMetricGroup,
TaskSlotTable<?> taskSlotTable,
java.util.function.Supplier<Long> managedMemoryTotalSupplier) |
static void |
MetricUtils.instantiateStatusMetrics(MetricGroup metricGroup) |
static void |
SystemResourcesMetricsInitializer.instantiateSystemMetrics(MetricGroup metricGroup,
Time probeInterval) |
Modifier and Type | Method and Description |
---|---|
DistributedRuntimeUDFContext |
BatchTask.createRuntimeContext(MetricGroup metrics) |
Constructor and Description |
---|
DistributedRuntimeUDFContext(TaskInfo taskInfo,
UserCodeClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
MetricGroup metrics,
ExternalResourceInfoProvider externalResourceInfoProvider) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
ShuffleIOOwnerContext.getInputGroup() |
MetricGroup |
ShuffleIOOwnerContext.getOutputGroup() |
MetricGroup |
ShuffleIOOwnerContext.getParentGroup() |
MetricGroup |
ShuffleEnvironmentContext.getParentMetricGroup() |
Modifier and Type | Method and Description |
---|---|
ShuffleIOOwnerContext |
ShuffleEnvironment.createShuffleIOOwnerContext(String ownerName,
ExecutionAttemptID executionAttemptID,
MetricGroup parentGroup)
Create a context of the shuffle input/output owner used to create partitions or gates
belonging to the owner.
|
Constructor and Description |
---|
ShuffleEnvironmentContext(Configuration configuration,
ResourceID taskExecutorResourceId,
MemorySize networkMemorySize,
boolean localCommunicationOnly,
InetAddress hostAddress,
TaskEventPublisher eventPublisher,
MetricGroup parentMetricGroup,
Executor ioExecutor) |
ShuffleIOOwnerContext(String ownerName,
ExecutionAttemptID executionAttemptID,
MetricGroup parentGroup,
MetricGroup outputGroup,
MetricGroup inputGroup) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
SourceCoordinatorContext.metricGroup() |
Modifier and Type | Method and Description |
---|---|
abstract <K> AbstractKeyedStateBackend<K> |
AbstractStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
<K> CheckpointableKeyedStateBackend<K> |
StateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry)
Creates a new
CheckpointableKeyedStateBackend that is responsible for holding
keyed state and checkpointing it. |
abstract <K> AbstractKeyedStateBackend<K> |
AbstractManagedMemoryStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry,
double managedMemoryFraction) |
default <K> CheckpointableKeyedStateBackend<K> |
StateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry,
double managedMemoryFraction)
Creates a new
CheckpointableKeyedStateBackend with the given managed memory fraction. |
Modifier and Type | Method and Description |
---|---|
<K> AbstractKeyedStateBackend<K> |
FsStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Modifier and Type | Method and Description |
---|---|
<K> AbstractKeyedStateBackend<K> |
MemoryStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Modifier and Type | Method and Description |
---|---|
static TaskManagerServices |
TaskManagerServices.fromConfiguration(TaskManagerServicesConfiguration taskManagerServicesConfiguration,
PermanentBlobService permanentBlobService,
MetricGroup taskManagerMetricGroup,
ExecutorService ioExecutor,
FatalErrorHandler fatalErrorHandler)
Creates and returns the task manager services.
|
Modifier and Type | Method and Description |
---|---|
MetricGroup |
StateBootstrapWrapperOperator.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
SavepointRuntimeContext.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
StreamOperator.getMetricGroup() |
MetricGroup |
AbstractStreamOperator.getMetricGroup() |
MetricGroup |
AbstractStreamOperatorV2.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
protected <K> CheckpointableKeyedStateBackend<K> |
StreamTaskStateInitializerImpl.keyedStatedBackend(TypeSerializer<K> keySerializer,
String operatorIdentifierText,
PrioritizedOperatorSubtaskState prioritizedOperatorSubtaskStates,
CloseableRegistry backendCloseableRegistry,
MetricGroup metricGroup,
double managedMemoryFraction) |
StreamOperatorStateContext |
StreamTaskStateInitializer.streamOperatorStateContext(OperatorID operatorID,
String operatorClassName,
ProcessingTimeService processingTimeService,
KeyContext keyContext,
TypeSerializer<?> keySerializer,
CloseableRegistry streamTaskCloseableRegistry,
MetricGroup metricGroup,
double managedMemoryFraction,
boolean isUsingCustomRawKeyedState)
Returns the
StreamOperatorStateContext for an AbstractStreamOperator that
runs in the stream task that owns this manager. |
StreamOperatorStateContext |
StreamTaskStateInitializerImpl.streamOperatorStateContext(OperatorID operatorID,
String operatorClassName,
ProcessingTimeService processingTimeService,
KeyContext keyContext,
TypeSerializer<?> keySerializer,
CloseableRegistry streamTaskCloseableRegistry,
MetricGroup metricGroup,
double managedMemoryFraction,
boolean isUsingCustomRawKeyedState) |
Constructor and Description |
---|
StreamingRuntimeContext(Environment env,
Map<String,Accumulator<?,?>> accumulators,
MetricGroup operatorMetricGroup,
OperatorID operatorID,
ProcessingTimeService processingTimeService,
KeyedStateStore keyedStateStore,
ExternalResourceInfoProvider externalResourceInfoProvider) |
Modifier and Type | Method and Description |
---|---|
<K> CheckpointableKeyedStateBackend<K> |
BatchExecutionStateBackend.createKeyedStateBackend(Environment env,
JobID jobID,
String operatorIdentifier,
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange,
TaskKvStateRegistry kvStateRegistry,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
TimestampsAndWatermarksContext.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
static <E> TimestampsAndWatermarks<E> |
TimestampsAndWatermarks.createNoOpEventTimeLogic(WatermarkStrategy<E> watermarkStrategy,
MetricGroup metrics) |
static <E> TimestampsAndWatermarks<E> |
TimestampsAndWatermarks.createProgressiveEventTimeLogic(WatermarkStrategy<E> watermarkStrategy,
MetricGroup metrics,
ProcessingTimeService timeService,
long periodicWatermarkIntervalMillis) |
Constructor and Description |
---|
TimestampsAndWatermarksContext(MetricGroup metricGroup) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
Evictor.EvictorContext.getMetricGroup()
Returns the metric group for this
Evictor . |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
Trigger.TriggerContext.getMetricGroup()
Returns the metric group for this
Trigger . |
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> subscribedPartitionsToStartOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup kafkaMetricGroup,
boolean useMetrics)
Creates the fetcher that connect to the Kafka brokers, pulls data, deserialized the data, and
emits it into the data streams.
|
Constructor and Description |
---|
AbstractFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> seedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
KafkaConsumerThread(org.slf4j.Logger log,
Handover handover,
Properties kafkaProperties,
ClosableBlockingQueue<KafkaTopicPartitionState<T,org.apache.kafka.common.TopicPartition>> unassignedPartitionsQueue,
String threadName,
long pollTimeout,
boolean useMetrics,
MetricGroup consumerMetricGroup,
MetricGroup subtaskMetricGroup) |
KafkaFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
KafkaShuffleFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics,
TypeSerializer<T> typeSerializer,
int producerParallelism) |
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaShuffleConsumer.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
Modifier and Type | Method and Description |
---|---|
protected RecordPublisher |
KinesisDataFetcher.createRecordPublisher(SequenceNumber sequenceNumber,
Properties configProps,
MetricGroup metricGroup,
StreamShardHandle subscribedShard) |
protected RecordPublisher |
DynamoDBStreamsDataFetcher.createRecordPublisher(SequenceNumber sequenceNumber,
Properties configProps,
MetricGroup metricGroup,
StreamShardHandle subscribedShard) |
protected ShardConsumer<T> |
KinesisDataFetcher.createShardConsumer(Integer subscribedShardStateIndex,
StreamShardHandle subscribedShard,
SequenceNumber lastSequenceNum,
MetricGroup metricGroup,
KinesisDeserializationSchema<T> shardDeserializer)
Create a new shard consumer.
|
Modifier and Type | Method and Description |
---|---|
RecordPublisher |
RecordPublisherFactory.create(StartingPosition startingPosition,
Properties consumerConfig,
MetricGroup metricGroup,
StreamShardHandle streamShardHandle)
Create a
RecordPublisher . |
Modifier and Type | Method and Description |
---|---|
FanOutRecordPublisher |
FanOutRecordPublisherFactory.create(StartingPosition startingPosition,
Properties consumerConfig,
MetricGroup metricGroup,
StreamShardHandle streamShardHandle)
Create a
FanOutRecordPublisher . |
Modifier and Type | Method and Description |
---|---|
PollingRecordPublisher |
PollingRecordPublisherFactory.create(StartingPosition startingPosition,
Properties consumerConfig,
MetricGroup metricGroup,
StreamShardHandle streamShardHandle)
Create a
PollingRecordPublisher . |
Constructor and Description |
---|
PollingRecordPublisherMetricsReporter(MetricGroup metricGroup) |
ShardConsumerMetricsReporter(MetricGroup metricGroup) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
WindowOperator.Context.getMetricGroup() |
Constructor and Description |
---|
LatencyStats(MetricGroup metricGroup,
int historySize,
int subtaskIndex,
OperatorID operatorID,
LatencyStats.Granularity granularity) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
FunctionContext.getMetricGroup()
Returns the metric group for this parallel subtask.
|
MetricGroup |
ConstantFunctionContext.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
Trigger.TriggerContext.getMetricGroup()
Returns the metric group for this
Trigger . |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.