From b7b2ef646c92b83844f0b9123f28a690eb9518cd Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 9 Apr 2021 10:42:12 +0300 Subject: [PATCH] Unregister JettyStatisticsCollector to fix memory leak in Broker shutdown - Unregister JettyStatisticsCollector from Prometheus client's default Collector registry at shutdown - Clear metricsServlet reference at shutdown of PulsarService --- .../org/apache/pulsar/broker/PulsarService.java | 2 ++ .../apache/pulsar/broker/web/WebService.java | 17 ++++++++++++++++- 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 7235aef7366ac..6e404038d4cdf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -339,6 +339,8 @@ public CompletableFuture closeAsync() { } } + metricsServlet = null; + if (this.webSocketService != null) { this.webSocketService.close(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java index 7f9c5808ddbc7..5d71176b22f18 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java @@ -19,6 +19,7 @@ package org.apache.pulsar.broker.web; import com.google.common.collect.Lists; +import io.prometheus.client.CollectorRegistry; import io.prometheus.client.jetty.JettyStatisticsCollector; import java.util.ArrayList; import java.util.EnumSet; @@ -72,6 +73,7 @@ public class WebService implements AutoCloseable { private final ServerConnector httpConnector; private final ServerConnector httpsConnector; + private JettyStatisticsCollector jettyStatisticsCollector; public WebService(PulsarService pulsar) throws PulsarServerException { this.handlers = Lists.newArrayList(); @@ -221,7 +223,8 @@ public void start() throws PulsarServerException { StatisticsHandler stats = new StatisticsHandler(); stats.setHandler(handlerCollection); try { - new JettyStatisticsCollector(stats).register(); + jettyStatisticsCollector = new JettyStatisticsCollector(stats); + jettyStatisticsCollector.register(); } catch (IllegalArgumentException e) { // Already registered. Eg: in unit tests } @@ -253,6 +256,18 @@ public void start() throws PulsarServerException { public void close() throws PulsarServerException { try { server.stop(); + // unregister statistics from Prometheus client's default CollectorRegistry singleton + // to prevent memory leaks in tests + if (jettyStatisticsCollector != null) { + try { + CollectorRegistry.defaultRegistry.unregister(jettyStatisticsCollector); + } catch (Exception e) { + // ignore any exception happening in unregister + // exception will be thrown for 2. instance of WebService in tests since + // the register supports a single JettyStatisticsCollector + } + jettyStatisticsCollector = null; + } webServiceExecutor.join(); log.info("Web service closed"); } catch (Exception e) {