From 7a184256861d4d2d24bc77f5689c6d2c386e732b Mon Sep 17 00:00:00 2001 From: maheshnikam <55378196+nikam14@users.noreply.github.com> Date: Sat, 23 Mar 2024 02:52:17 +0530 Subject: [PATCH 1/4] Update ReplicatedSubscriptionsController.java --- .../persistent/ReplicatedSubscriptionsController.java | 6 ++++++ 1 file changed, 6 insertions(+) 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 e011ed8d660f6..499198bf6a75a 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; @@ -75,6 +76,10 @@ public class ReplicatedSubscriptionsController implements AutoCloseable, Topic.P "Counter of currently pending snapshots") .register(); + private static final Counter timedoutSnapshotsMetric = Counter + .build().name("pulsar_replicated_subscriptions_snapshot_timeouts") + .help("Counter of timed out snapshots").register(); + public ReplicatedSubscriptionsController(PersistentTopic topic, String localCluster) { this.topic = topic; this.localCluster = localCluster; @@ -257,6 +262,7 @@ private void cleanupTimedOutSnapshots() { } pendingSnapshotsMetric.dec(); + timedoutSnapshotsMetric.inc(); it.remove(); } } From 48fe420cc834dd26599eb3e31f2fcc350fffbfeb Mon Sep 17 00:00:00 2001 From: maheshnikam <55378196+nikam14@users.noreply.github.com> Date: Sat, 23 Mar 2024 11:57:45 +0530 Subject: [PATCH 2/4] Update ReplicatedSubscriptionsController.java --- .../service/persistent/ReplicatedSubscriptionsController.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 499198bf6a75a..941f911a966a0 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 @@ -79,7 +79,7 @@ public class ReplicatedSubscriptionsController implements AutoCloseable, Topic.P private static final Counter timedoutSnapshotsMetric = Counter .build().name("pulsar_replicated_subscriptions_snapshot_timeouts") .help("Counter of timed out snapshots").register(); - + public ReplicatedSubscriptionsController(PersistentTopic topic, String localCluster) { this.topic = topic; this.localCluster = localCluster; From 951581983ef259ad26adda0e8371856c5c0872ea Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 10 Oct 2024 15:10:52 +0300 Subject: [PATCH 3/4] Add annotation for PIP-264 / OTel --- .../persistent/ReplicatedSubscriptionsController.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) 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 49d664c05101b..8d4795ff0b25a 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 @@ -81,13 +81,16 @@ 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_snapshot_timeouts") .help("Counter of timed out snapshots").register(); private final OpenTelemetryReplicatedSubscriptionStats stats; - public ReplicatedSubscriptionsController(PersistentTopic topic, String localCluster) { this.topic = topic; this.localCluster = localCluster; From 5b6968577078009616d89eb1188dda83a57b7c9b Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 10 Oct 2024 15:17:34 +0300 Subject: [PATCH 4/4] Revisit metric name to make it consistent with the other metric name --- .../service/persistent/ReplicatedSubscriptionsController.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 8d4795ff0b25a..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 @@ -86,7 +86,7 @@ public class ReplicatedSubscriptionsController implements AutoCloseable, Topic.P newMetricName = OpenTelemetryReplicatedSubscriptionStats.SNAPSHOT_OPERATION_COUNT_METRIC_NAME) @Deprecated private static final Counter timedoutSnapshotsMetric = Counter - .build().name("pulsar_replicated_subscriptions_snapshot_timeouts") + .build().name("pulsar_replicated_subscriptions_timedout_snapshots") .help("Counter of timed out snapshots").register(); private final OpenTelemetryReplicatedSubscriptionStats stats;