diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/ReplicatedSubscriptionsController.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/ReplicatedSubscriptionsController.java index f56cf9de66b75..4fb0022194a02 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/ReplicatedSubscriptionsController.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/ReplicatedSubscriptionsController.java @@ -20,6 +20,7 @@ import static org.apache.pulsar.common.util.Runnables.catchingAndLoggingThrowables; import io.netty.buffer.ByteBuf; +import io.prometheus.client.Counter; import io.prometheus.client.Gauge; import java.io.IOException; import java.time.Clock; @@ -80,6 +81,14 @@ public class ReplicatedSubscriptionsController implements AutoCloseable, Topic.P "Counter of currently pending snapshots") .register(); + // timeouts use SnapshotOperationResult.TIMEOUT.attributes on the same metric + @PulsarDeprecatedMetric( + newMetricName = OpenTelemetryReplicatedSubscriptionStats.SNAPSHOT_OPERATION_COUNT_METRIC_NAME) + @Deprecated + private static final Counter timedoutSnapshotsMetric = Counter + .build().name("pulsar_replicated_subscriptions_timedout_snapshots") + .help("Counter of timed out snapshots").register(); + private final OpenTelemetryReplicatedSubscriptionStats stats; public ReplicatedSubscriptionsController(PersistentTopic topic, String localCluster) { @@ -263,6 +272,7 @@ private void cleanupTimedOutSnapshots() { } pendingSnapshotsMetric.dec(); + timedoutSnapshotsMetric.inc(); var latencyMillis = entry.getValue().getDurationMillis(); stats.recordSnapshotTimedOut(latencyMillis); it.remove();