Modifier and Type | Method and Description |
---|---|
MetricGroup |
WatermarkGeneratorSupplier.Context.getMetricGroup()
Returns the metric group for the context in which the created
WatermarkGenerator
is used. |
MetricGroup |
TimestampAssignerSupplier.Context.getMetricGroup()
Returns the metric group for the context in which the created
TimestampAssigner
is used. |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
SerializationSchema.InitializationContext.getMetricGroup()
Returns the metric group for the parallel subtask of the source that runs this
SerializationSchema . |
MetricGroup |
DeserializationSchema.InitializationContext.getMetricGroup()
Returns the metric group for the parallel subtask of the source that runs this
DeserializationSchema . |
Modifier and Type | Class and Description |
---|---|
class |
ChangelogStorageMetricGroup
Metrics related to the Changelog Storage used by the Changelog State Backend.
|
Constructor and Description |
---|
ChangelogStorageMetricGroup(MetricGroup parent) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
PulsarDeserializationSchemaInitializationContext.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
DummyInitializationContext.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
TestingDeserializationContext.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
<K> AbstractKeyedStateBackend<K> |
EmbeddedRocksDBStateBackend.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)
Deprecated.
|
<K> AbstractKeyedStateBackend<K> |
EmbeddedRocksDBStateBackend.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) |
<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)
Deprecated.
|
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,
EmbeddedRocksDBStateBackend.PriorityQueueStateType priorityQueueStateType,
TtlTimeProvider ttlTimeProvider,
LatencyTrackingStateConfig latencyTrackingStateConfig,
MetricGroup metricGroup,
Collection<KeyedStateHandle> stateHandles,
StreamCompressionDecorator keyGroupCompressionDecorator,
CloseableRegistry cancelStreamRegistry) |
RocksDBNativeMetricMonitor(RocksDBNativeMetricOptions options,
MetricGroup metricGroup,
org.rocksdb.RocksDB rocksDB,
org.rocksdb.Statistics statistics) |
Constructor and Description |
---|
RocksDBFullRestoreOperation(KeyGroupRange keyGroupRange,
ClassLoader userCodeClassLoader,
Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
StateSerializerProvider<K> keySerializerProvider,
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) |
RocksDBHeapTimersFullRestoreOperation(KeyGroupRange keyGroupRange,
int numberOfKeyGroups,
ClassLoader userCodeClassLoader,
Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
LinkedHashMap<String,HeapPriorityQueueSnapshotRestoreWrapper<?>> registeredPQStates,
HeapPriorityQueueSetFactory priorityQueueFactory,
StateSerializerProvider<K> keySerializerProvider,
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) |
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,
double overlapFractionThreshold) |
RocksDBNoneRestoreOperation(Map<String,RocksDBKeyedStateBackend.RocksDbKvStateInfo> kvStateInformation,
File instanceRocksDBPath,
org.rocksdb.DBOptions dbOptions,
java.util.function.Function<String,org.rocksdb.ColumnFamilyOptions> columnFamilyOptionsFactory,
RocksDBNativeMetricOptions nativeMetricOptions,
MetricGroup metricGroup,
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 |
---|---|
default 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.
|
MetricGroup |
LogicalScopeProvider.getWrappedMetricGroup()
Returns the underlying metric group.
|
Modifier and Type | Method and Description |
---|---|
static LogicalScopeProvider |
LogicalScopeProvider.castFrom(MetricGroup metricGroup)
Casts the given metric group to a
LogicalScopeProvider , if it implements the
interface. |
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 | Interface and Description |
---|---|
interface |
CacheMetricGroup
Pre-defined metrics for cache.
|
interface |
OperatorCoordinatorMetricGroup
Special
MetricGroup representing an Operator coordinator. |
interface |
OperatorIOMetricGroup
Metric group that contains shareable pre-defined IO-related metrics for operators.
|
interface |
OperatorMetricGroup
Special
MetricGroup representing an Operator. |
interface |
SinkWriterMetricGroup
Pre-defined metrics for sinks.
|
interface |
SourceReaderMetricGroup
Pre-defined metrics for
SourceReader . |
interface |
SplitEnumeratorMetricGroup
Pre-defined metrics for
SplitEnumerator . |
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(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 |
AbstractReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
MetricReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group)
Called when a new
Metric was added. |
void |
AbstractReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
MetricReporter.notifyOfRemovedMetric(Metric metric,
String metricName,
MetricGroup group)
Called when a
Metric was removed. |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
MetricListener.getMetricGroup()
Get the root metric group of this listener.
|
Constructor and Description |
---|
FlinkMetricContainer(MetricGroup metricGroup) |
Constructor and Description |
---|
CheckpointStatsTracker(int numRememberedCheckpoints,
MetricGroup metricGroup)
Creates a new checkpoint stats tracker.
|
Modifier and Type | Method and Description |
---|---|
static NettyShuffleEnvironment |
NettyShuffleServiceFactory.createNettyShuffleEnvironment(NettyShuffleEnvironmentConfiguration config,
ResourceID taskExecutorResourceId,
TaskEventPublisher taskEventPublisher,
ResultPartitionManager resultPartitionManager,
ConnectionManager connectionManager,
MetricGroup metricGroup,
Executor ioExecutor,
int numberOfSlots,
String[] tmpDirPaths) |
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.registerDebloatingTaskMetrics(SingleInputGate[] inputGates,
MetricGroup taskGroup) |
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 | 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 |
InternalCacheMetricGroup
A
CacheMetricGroup which register all cache related metrics under a subgroup of the
parent metric group. |
class |
InternalOperatorIOMetricGroup
Metric group that contains shareable pre-defined IO-related metrics.
|
class |
InternalOperatorMetricGroup
Special
MetricGroup representing an Operator. |
class |
InternalSinkWriterMetricGroup
Special
MetricGroup representing an Operator. |
class |
InternalSourceReaderMetricGroup
Special
MetricGroup representing an Operator. |
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 |
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
InternalOperatorMetricGroup 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 |
AbstractMetricGroup.addGroup(String name) |
MetricGroup |
ProxyMetricGroup.addGroup(String name) |
MetricGroup |
AbstractMetricGroup.addGroup(String key,
String value) |
MetricGroup |
ProxyMetricGroup.addGroup(String key,
String value) |
MetricGroup |
GenericKeyMetricGroup.addGroup(String key,
String value) |
MetricGroup |
InternalOperatorMetricGroup.getJobMetricGroup() |
MetricGroup |
FrontMetricGroup.getWrappedMetricGroup() |
Modifier and Type | Method and Description |
---|---|
static InternalSourceReaderMetricGroup |
InternalSourceReaderMetricGroup.mock(MetricGroup metricGroup) |
static InternalSinkWriterMetricGroup |
InternalSinkWriterMetricGroup.mock(MetricGroup metricGroup) |
static InternalSinkWriterMetricGroup |
InternalSinkWriterMetricGroup.mock(MetricGroup metricGroup,
OperatorIOMetricGroup operatorIOMetricGroup) |
Constructor and Description |
---|
InternalCacheMetricGroup(MetricGroup parentMetricGroup,
String subGroupName)
Creates a subgroup with the specified subgroup name under the parent group.
|
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 |
---|---|
static void |
SchedulerBase.registerJobMetrics(MetricGroup metrics,
JobStatusProvider jobStatusProvider,
Gauge<Long> numberOfRestarts,
DeploymentStateTimeMetrics deploymentTimeMetrics,
java.util.function.Consumer<JobStatusListener> jobStatusListenerRegistrar,
long initializationTimestamp,
MetricOptions.JobStatusMetricsSettings jobStatusMetricsSettings) |
Modifier and Type | Method and Description |
---|---|
static void |
StateTimeMetric.register(MetricOptions.JobStatusMetricsSettings jobStatusMetricsSettings,
MetricGroup metricGroup,
StateTimeMetric stateTimeMetric,
String baseName) |
void |
MetricsRegistrar.registerMetrics(MetricGroup metricGroup) |
void |
DeploymentStateTimeMetrics.registerMetrics(MetricGroup metricGroup) |
void |
JobStatusMetrics.registerMetrics(MetricGroup metricGroup) |
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,
int numberOfSlots,
String[] tmpDirPaths,
TaskEventPublisher eventPublisher,
MetricGroup parentMetricGroup,
Executor ioExecutor) |
ShuffleIOOwnerContext(String ownerName,
ExecutionAttemptID executionAttemptID,
MetricGroup parentGroup,
MetricGroup outputGroup,
MetricGroup inputGroup) |
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. |
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. |
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) |
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)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
<K> AbstractKeyedStateBackend<K> |
HashMapStateBackend.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)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
MetricGroup |
LatencyTrackingStateConfig.getMetricGroup() |
Modifier and Type | Method and Description |
---|---|
LatencyTrackingStateConfig.Builder |
LatencyTrackingStateConfig.Builder.setMetricGroup(MetricGroup metricGroup) |
Modifier and Type | Method and Description |
---|---|
static TaskManagerServices |
TaskManagerServices.fromConfiguration(TaskManagerServicesConfiguration taskManagerServicesConfiguration,
PermanentBlobService permanentBlobService,
MetricGroup taskManagerMetricGroup,
ExecutorService ioExecutor,
FatalErrorHandler fatalErrorHandler,
WorkingDirectory workingDirectory)
Creates and returns the task manager services.
|
Modifier and Type | Method and Description |
---|---|
<K> CheckpointableKeyedStateBackend<K> |
AbstractChangelogStateBackend.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> |
AbstractChangelogStateBackend.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) |
protected <K> CheckpointableKeyedStateBackend<K> |
DeactivatedChangelogStateBackend.restore(Environment env,
String operatorIdentifier,
KeyGroupRange keyGroupRange,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<ChangelogStateBackendHandle> stateBackendHandles,
ChangelogBackendRestoreOperation.BaseBackendBuilder<K> baseBackendBuilder) |
protected <K> CheckpointableKeyedStateBackend<K> |
ChangelogStateBackend.restore(Environment env,
String operatorIdentifier,
KeyGroupRange keyGroupRange,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<ChangelogStateBackendHandle> stateBackendHandles,
ChangelogBackendRestoreOperation.BaseBackendBuilder<K> baseBackendBuilder) |
protected abstract <K> CheckpointableKeyedStateBackend<K> |
AbstractChangelogStateBackend.restore(Environment env,
String operatorIdentifier,
KeyGroupRange keyGroupRange,
TtlTimeProvider ttlTimeProvider,
MetricGroup metricGroup,
Collection<ChangelogStateBackendHandle> stateBackendHandles,
ChangelogBackendRestoreOperation.BaseBackendBuilder<K> baseBackendBuilder) |
Modifier and Type | Class and Description |
---|---|
class |
ChangelogMaterializationMetricGroup
Metrics related to the materialization part of Changelog.
|
Constructor and Description |
---|
ChangelogMaterializationMetricGroup(MetricGroup parentMetricGroup) |
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 |
StreamTaskStateInitializerImpl.streamOperatorStateContext(OperatorID operatorID,
String operatorClassName,
ProcessingTimeService processingTimeService,
KeyContext keyContext,
TypeSerializer<?> keySerializer,
CloseableRegistry streamTaskCloseableRegistry,
MetricGroup metricGroup,
double managedMemoryFraction,
boolean isUsingCustomRawKeyedState) |
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. |
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)
Deprecated.
|
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 |
DynamoDBStreamsDataFetcher.createRecordPublisher(SequenceNumber sequenceNumber,
Properties configProps,
MetricGroup metricGroup,
StreamShardHandle subscribedShard) |
protected RecordPublisher |
KinesisDataFetcher.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.
|
Modifier and Type | Method and Description |
---|---|
MetricGroup |
Trigger.TriggerContext.getMetricGroup()
Returns the metric group for this
Trigger . |
Copyright © 2014–2023 The Apache Software Foundation. All rights reserved.