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,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Accumulator<?,?>> accumulators,
Map<String,Future<Path>> cpTasks,
MetricGroup metrics) |
RuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
MetricGroup metrics) |
Modifier and Type | Method and Description |
---|---|
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.
|
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) |
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 |
PrometheusReporter.notifyOfAddedMetric(Metric metric,
String metricName,
MetricGroup group) |
void |
PrometheusReporter.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 should be removed. |
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 timeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
int parallelismForAutoMax,
BlobWriter blobWriter,
org.slf4j.Logger log)
Builds the ExecutionGraph from the JobGraph.
|
Modifier and Type | Method and Description |
---|---|
void |
RestartIndividualStrategy.registerMetrics(MetricGroup metricGroup) |
void |
FailoverStrategy.registerMetrics(MetricGroup metricGroup)
Tells the FailoverStrategy to register its metrics.
|
Modifier and Type | Method and Description |
---|---|
static void |
ResultPartitionMetrics.registerQueueLengthMetrics(MetricGroup group,
ResultPartition partition) |
Modifier and Type | Method and Description |
---|---|
static void |
InputGateMetrics.registerQueueLengthMetrics(MetricGroup group,
SingleInputGate gate) |
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 |
GenericMetricGroup
A simple named
MetricGroup that is used to hold
subgroups of metrics. |
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 |
ProxyMetricGroup<P extends MetricGroup>
Metric group which forwards all registration calls to its parent metric group.
|
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. |
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) |
Modifier and Type | Method and Description |
---|---|
static void |
MetricUtils.instantiateStatusMetrics(MetricGroup metricGroup) |
Modifier and Type | Method and Description |
---|---|
DistributedRuntimeUDFContext |
BatchTask.createRuntimeContext(MetricGroup metrics) |
Constructor and Description |
---|
DistributedRuntimeUDFContext(TaskInfo taskInfo,
ClassLoader userCodeClassLoader,
ExecutionConfig executionConfig,
Map<String,Future<Path>> cpTasks,
Map<String,Accumulator<?,?>> accumulators,
MetricGroup metrics) |
Modifier and Type | Field and Description |
---|---|
protected MetricGroup |
AbstractStreamOperator.metrics
Metric group for the operator.
|
Modifier and Type | Method and Description |
---|---|
MetricGroup |
StreamOperator.getMetricGroup() |
MetricGroup |
AbstractStreamOperator.getMetricGroup() |
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,?> |
FlinkKafkaConsumer010.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer09.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer08.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> subscribedPartitionsToStartOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
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 |
---|
Kafka010Fetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
Kafka09Fetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
KafkaConsumerThread(org.slf4j.Logger log,
Handover handover,
Properties kafkaProperties,
ClosableBlockingQueue<KafkaTopicPartitionState<org.apache.kafka.common.TopicPartition>> unassignedPartitionsQueue,
KafkaConsumerCallBridge consumerCallBridge,
String threadName,
long pollTimeout,
boolean useMetrics,
MetricGroup subtaskMetricGroup) |
Constructor and Description |
---|
AbstractFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> seedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
Kafka08Fetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> seedPartitionsWithInitialOffsets,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long autoCommitInterval,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
Modifier and Type | Method and Description |
---|---|
MetricGroup |
WindowOperator.Context.getMetricGroup() |
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.