From becdde7a13dbd55db1ee9349ee9dc7c06700fe2c Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sat, 25 Apr 2020 10:39:58 +0800 Subject: [PATCH 01/12] expose managed ledger bookie client metric to prometheus --- conf/broker.conf | 6 + conf/standalone.conf | 6 + .../mledger/ManagedLedgerFactory.java | 7 + .../mledger/ManagedLedgerFactoryConfig.java | 10 + .../impl/ManagedLedgerFactoryImpl.java | 24 +++ .../prometheus/DataSketchesOpStatsLogger.java | 201 ++++++++++++++++++ .../stats/prometheus/LongAdderCounter.java | 58 +++++ .../prometheus/PrometheusMetricsProvider.java | 154 ++++++++++++++ .../prometheus/PrometheusStatsLogger.java | 77 +++++++ .../prometheus/PrometheusTextFormatUtil.java | 152 +++++++++++++ .../mledger/stats/prometheus/SimpleGauge.java | 37 ++++ .../pulsar/broker/ServiceConfiguration.java | 13 ++ .../broker/ManagedLedgerClientFactory.java | 2 + .../PrometheusMetricsGenerator.java | 19 ++ 14 files changed, 766 insertions(+) create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java diff --git a/conf/broker.conf b/conf/broker.conf index 8aebf1038716f..bc49086acb8d5 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -680,6 +680,12 @@ managedLedgerReadEntryTimeoutSeconds=0 # Add entry timeout when broker tries to publish message to bookkeeper (0 to disable it). managedLedgerAddEntryTimeoutSeconds=0 +# Managed ledger prometheus stats latency rollover seconds (default: 60s) +managedLedgerPrometheusStatsLatencyRolloverSeconds=60 + +# Whether trace managed ledger task execution time +managedLedgerTraceTaskExecution=false + ### --- Load balancer --- ### # Enable load balancer diff --git a/conf/standalone.conf b/conf/standalone.conf index 43138f8677bce..28cf754bf81df 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -449,6 +449,12 @@ managedLedgerAddEntryTimeoutSeconds=0 # Use Open Range-Set to cache unacked messages managedLedgerUnackedRangesOpenCacheSetEnabled=true +# Managed ledger prometheus stats latency rollover seconds (default: 60s) +managedLedgerPrometheusStatsLatencyRolloverSeconds=60 + +# Whether trace managed ledger task execution time +managedLedgerTraceTaskExecution=false + ### --- Load balancer --- ### loadManagerClassName=org.apache.pulsar.broker.loadbalance.NoopLoadManager diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java index a28249af4b82a..286b5f72761b2 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java @@ -25,6 +25,7 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.ManagedLedgerInfoCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenLedgerCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.OpenReadOnlyCursorCallback; +import org.apache.bookkeeper.stats.StatsProvider; /** * A factory to open/create managed ledgers and delete them. @@ -140,4 +141,10 @@ void asyncOpenReadOnlyCursor(String managedLedgerName, Position startPosition, M */ void shutdown() throws InterruptedException, ManagedLedgerException; + /** + * Get managed ledger stats provider. + * + * @return StatsProvider + */ + StatsProvider getStatsProvider(); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index 4d1eb3121af12..d8e405e09bf47 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -51,4 +51,14 @@ public class ManagedLedgerFactoryConfig { * Whether we should make a copy of the entry payloads when inserting in cache */ private boolean copyEntriesInCache = false; + + /** + * Whether trace managed ledger task execution time + */ + private boolean traceTaskExecution = false; + + /** + * managed ledger prometheus stats Latency Rollover Seconds + */ + private int prometheusStatsLatencyRolloverSeconds = 60; } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 66ad3baf0085a..a08366e0cb55d 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -73,8 +73,13 @@ import org.apache.bookkeeper.mledger.proto.MLDataFormats.LongProperty; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.MessageRange; +import org.apache.bookkeeper.mledger.stats.prometheus.PrometheusMetricsProvider; import org.apache.bookkeeper.mledger.util.Futures; +import org.apache.bookkeeper.stats.NullStatsLogger; +import org.apache.bookkeeper.stats.StatsLogger; +import org.apache.bookkeeper.stats.StatsProvider; import org.apache.bookkeeper.zookeeper.ZooKeeperClient; +import org.apache.commons.configuration.Configuration; import org.apache.pulsar.common.util.DateFormatter; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.Stat; @@ -107,6 +112,9 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private static final int StatsPeriodSeconds = 60; + private StatsLogger statsLogger = NullStatsLogger.INSTANCE; + private StatsProvider statsProvider = new PrometheusMetricsProvider(); + public ManagedLedgerFactoryImpl(ClientConfiguration bkClientConfiguration, String zkConnection) throws Exception { this(bkClientConfiguration, zkConnection, new ManagedLedgerFactoryConfig()); } @@ -149,12 +157,22 @@ public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolic private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, boolean isBookkeeperManaged, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config) throws Exception { + Configuration configuration = new ClientConfiguration(); + configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS, config.getPrometheusStatsLatencyRolloverSeconds()); + statsProvider.start(configuration); + + statsLogger = statsProvider.getStatsLogger("pulsar_bookie_client"); + scheduledExecutor = OrderedScheduler.newSchedulerBuilder() .numThreads(config.getNumManagedLedgerSchedulerThreads()) + .statsLogger(statsLogger) + .traceTaskExecution(config.isTraceTaskExecution()) .name("bookkeeper-ml-scheduler") .build(); orderedExecutor = OrderedExecutor.newBuilder() .numThreads(config.getNumManagedLedgerWorkerThreads()) + .statsLogger(statsLogger) + .traceTaskExecution(config.isTraceTaskExecution()) .name("bookkeeper-ml-workers") .build(); cacheEvictionExecutor = Executors @@ -197,6 +215,11 @@ public BookKeeper get(EnsemblePlacementPolicyConfig policy) { } } + @Override + public StatsProvider getStatsProvider() { + return statsProvider; + } + private synchronized void refreshStats() { long now = System.nanoTime(); long period = now - lastStatTimestamp; @@ -462,6 +485,7 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) { scheduledExecutor.shutdownNow(); orderedExecutor.shutdownNow(); cacheEvictionExecutor.shutdownNow(); + statsProvider.stop(); entryCacheManager.clear(); try { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java new file mode 100644 index 0000000000000..dc5cdd76993df --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java @@ -0,0 +1,201 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.stats.prometheus; + +import com.yahoo.sketches.quantiles.DoublesSketch; +import com.yahoo.sketches.quantiles.DoublesSketchBuilder; +import com.yahoo.sketches.quantiles.DoublesUnion; +import com.yahoo.sketches.quantiles.DoublesUnionBuilder; + +import io.netty.util.concurrent.FastThreadLocal; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.LongAdder; +import java.util.concurrent.locks.StampedLock; + +import org.apache.bookkeeper.stats.OpStatsData; +import org.apache.bookkeeper.stats.OpStatsLogger; + +/** + * OpStatsLogger implementation that uses DataSketches library to calculate the approximated latency quantiles. + */ +public class DataSketchesOpStatsLogger implements OpStatsLogger { + + /** + * Use 2 rotating thread local accessor so that we can safely swap them. + */ + private volatile ThreadLocalAccessor current; + private volatile ThreadLocalAccessor replacement; + + /** + * These are the sketches where all the aggregated results are published. + */ + private volatile DoublesSketch successResult; + private volatile DoublesSketch failResult; + + private final LongAdder successCountAdder = new LongAdder(); + private final LongAdder failCountAdder = new LongAdder(); + + private final LongAdder successSumAdder = new LongAdder(); + private final LongAdder failSumAdder = new LongAdder(); + + DataSketchesOpStatsLogger() { + this.current = new ThreadLocalAccessor(); + this.replacement = new ThreadLocalAccessor(); + } + + @Override + public void registerFailedEvent(long eventLatency, TimeUnit unit) { + double valueMillis = unit.toMicros(eventLatency) / 1000.0; + + failCountAdder.increment(); + failSumAdder.add((long) valueMillis); + + LocalData localData = current.localData.get(); + + long stamp = localData.lock.readLock(); + try { + localData.failSketch.update(valueMillis); + } finally { + localData.lock.unlockRead(stamp); + } + } + + @Override + public void registerSuccessfulEvent(long eventLatency, TimeUnit unit) { + double valueMillis = unit.toMicros(eventLatency) / 1000.0; + + successCountAdder.increment(); + successSumAdder.add((long) valueMillis); + + LocalData localData = current.localData.get(); + + long stamp = localData.lock.readLock(); + try { + localData.successSketch.update(valueMillis); + } finally { + localData.lock.unlockRead(stamp); + } + } + + @Override + public void registerSuccessfulValue(long value) { + successCountAdder.increment(); + successSumAdder.add(value); + + LocalData localData = current.localData.get(); + + long stamp = localData.lock.readLock(); + try { + localData.successSketch.update(value); + } finally { + localData.lock.unlockRead(stamp); + } + } + + @Override + public void registerFailedValue(long value) { + failCountAdder.increment(); + failSumAdder.add(value); + + LocalData localData = current.localData.get(); + + long stamp = localData.lock.readLock(); + try { + localData.failSketch.update(value); + } finally { + localData.lock.unlockRead(stamp); + } + } + + @Override + public OpStatsData toOpStatsData() { + // Not relevant as we don't use JMX here + throw new UnsupportedOperationException(); + } + + @Override + public void clear() { + // Not relevant as we don't use JMX here + throw new UnsupportedOperationException(); + } + + public void rotateLatencyCollection() { + // Swap current with replacement + ThreadLocalAccessor local = current; + current = replacement; + replacement = local; + + final DoublesUnion aggregateSuccesss = new DoublesUnionBuilder().build(); + final DoublesUnion aggregateFail = new DoublesUnionBuilder().build(); + local.map.forEach((localData, b) -> { + long stamp = localData.lock.writeLock(); + try { + aggregateSuccesss.update(localData.successSketch); + localData.successSketch.reset(); + aggregateFail.update(localData.failSketch); + localData.failSketch.reset(); + } finally { + localData.lock.unlockWrite(stamp); + } + }); + + successResult = aggregateSuccesss.getResultAndReset(); + failResult = aggregateFail.getResultAndReset(); + } + + public long getCount(boolean success) { + return success ? successCountAdder.sum() : failCountAdder.sum(); + } + + public long getSum(boolean success) { + return success ? successSumAdder.sum() : failSumAdder.sum(); + } + + public double getQuantileValue(boolean success, double quantile) { + DoublesSketch s = success ? successResult : failResult; + return s != null ? s.getQuantile(quantile) : Double.NaN; + } + + private static class LocalData { + private final DoublesSketch successSketch = new DoublesSketchBuilder().build(); + private final DoublesSketch failSketch = new DoublesSketchBuilder().build(); + private final StampedLock lock = new StampedLock(); + } + + private static class ThreadLocalAccessor { + private final Map map = new ConcurrentHashMap<>(); + private final FastThreadLocal localData = new FastThreadLocal() { + + @Override + protected LocalData initialValue() throws Exception { + LocalData localData = new LocalData(); + map.put(localData, Boolean.TRUE); + return localData; + } + + @Override + protected void onRemoval(LocalData value) throws Exception { + map.remove(value); + } + }; + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java new file mode 100644 index 0000000000000..44fc4fafb6e2e --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java @@ -0,0 +1,58 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.stats.prometheus; + +import java.util.concurrent.atomic.LongAdder; + +import org.apache.bookkeeper.stats.Counter; + +/** + * {@link Counter} implementation based on {@link LongAdder}. + * + *

LongAdder keeps a counter per-thread and then aggregates to get the result, in order to avoid contention between + * multiple threads. + */ +public class LongAdderCounter implements Counter { + private final LongAdder counter = new LongAdder(); + + @Override + public void clear() { + counter.reset(); + } + + @Override + public void inc() { + counter.increment(); + } + + @Override + public void dec() { + counter.decrement(); + } + + @Override + public void add(long delta) { + counter.add(delta); + } + + @Override + public Long get() { + return counter.sum(); + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java new file mode 100644 index 0000000000000..b9c53bd9b37c6 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java @@ -0,0 +1,154 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.stats.prometheus; + +import com.google.common.annotations.VisibleForTesting; + +import io.netty.util.concurrent.DefaultThreadFactory; +import io.netty.util.internal.PlatformDependent; +import io.prometheus.client.Collector; +import io.prometheus.client.CollectorRegistry; + +import java.io.IOException; +import java.io.Writer; +import java.lang.reflect.Field; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicLong; + +import org.apache.bookkeeper.stats.CachingStatsProvider; +import org.apache.bookkeeper.stats.StatsLogger; +import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.commons.configuration.Configuration; +import org.apache.commons.lang.StringUtils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * A Prometheus based {@link StatsProvider} implementation. + */ +public class PrometheusMetricsProvider implements StatsProvider { + private ScheduledExecutorService executor; + + public static final String PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS = "prometheusStatsLatencyRolloverSeconds"; + public static final int DEFAULT_PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS = 60; + + final CollectorRegistry registry; + private final CachingStatsProvider cachingStatsProvider; + + /** + * These acts a registry of the metrics defined in this provider + */ + final ConcurrentMap counters = new ConcurrentSkipListMap<>(); + final ConcurrentMap> gauges = new ConcurrentSkipListMap<>(); + final ConcurrentMap opStats = new ConcurrentSkipListMap<>(); + + public PrometheusMetricsProvider() { + this(CollectorRegistry.defaultRegistry); + } + + public PrometheusMetricsProvider(CollectorRegistry registry) { + this.registry = registry; + this.cachingStatsProvider = new CachingStatsProvider(new StatsProvider() { + @Override + public void start(Configuration conf) { + // nop + } + + @Override + public void stop() { + // nop + } + + @Override + public StatsLogger getStatsLogger(String scope) { + return new PrometheusStatsLogger(PrometheusMetricsProvider.this, scope); + } + + @Override + public String getStatsName(String... statsComponents) { + String completeName; + if (statsComponents.length == 0) { + return ""; + } else if (statsComponents[0].isEmpty()) { + completeName = StringUtils.join(statsComponents, '_', 1, statsComponents.length); + } else { + completeName = StringUtils.join(statsComponents, '_'); + } + return Collector.sanitizeMetricName(completeName); + } + }); + } + + @Override + public void start(Configuration conf) { + executor = Executors.newSingleThreadScheduledExecutor(new DefaultThreadFactory("metrics")); + + int latencyRolloverSeconds = conf.getInt(PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS, + DEFAULT_PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS); + + executor.scheduleAtFixedRate(() -> { + rotateLatencyCollection(); + }, 1, latencyRolloverSeconds, TimeUnit.SECONDS); + } + + @Override + public void stop() { + executor.shutdownNow(); + } + + @Override + public StatsLogger getStatsLogger(String scope) { + return this.cachingStatsProvider.getStatsLogger(scope); + } + + @Override + public void writeAllMetrics(Writer writer) throws IOException { + PrometheusTextFormatUtil.writeMetricsCollectedByPrometheusClient(writer, registry); + + gauges.forEach((name, gauge) -> PrometheusTextFormatUtil.writeGauge(writer, name, gauge)); + counters.forEach((name, counter) -> PrometheusTextFormatUtil.writeCounter(writer, name, counter)); + opStats.forEach((name, opStatLogger) -> PrometheusTextFormatUtil.writeOpStat(writer, name, opStatLogger)); + } + + @Override + public String getStatsName(String... statsComponents) { + return cachingStatsProvider.getStatsName(statsComponents); + } + + @VisibleForTesting + void rotateLatencyCollection() { + opStats.forEach((name, metric) -> { + metric.rotateLatencyCollection(); + }); + } + + private void registerMetrics(Collector collector) { + try { + collector.register(registry); + } catch (Exception e) { + // Ignore if these were already registered + if (log.isDebugEnabled()) { + log.debug("Failed to register Prometheus collector exports", e); + } + } + } + + + private static final Logger log = LoggerFactory.getLogger(org.apache.bookkeeper.stats.prometheus.PrometheusMetricsProvider.class); +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java new file mode 100644 index 0000000000000..295d4575ca0ca --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java @@ -0,0 +1,77 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.stats.prometheus; + +import com.google.common.base.Joiner; + +import io.prometheus.client.Collector; + +import org.apache.bookkeeper.stats.Counter; +import org.apache.bookkeeper.stats.Gauge; +import org.apache.bookkeeper.stats.OpStatsLogger; +import org.apache.bookkeeper.stats.StatsLogger; + +/** + * A {@code Prometheus} based {@link StatsLogger} implementation. + */ +public class PrometheusStatsLogger implements StatsLogger { + + private final PrometheusMetricsProvider provider; + private final String scope; + + PrometheusStatsLogger(PrometheusMetricsProvider provider, String scope) { + this.provider = provider; + this.scope = scope; + } + + @Override + public OpStatsLogger getOpStatsLogger(String name) { + return provider.opStats.computeIfAbsent(completeName(name), x -> new DataSketchesOpStatsLogger()); + } + + @Override + public Counter getCounter(String name) { + return provider.counters.computeIfAbsent(completeName(name), x -> new LongAdderCounter()); + } + + @Override + public void registerGauge(String name, Gauge gauge) { + provider.gauges.computeIfAbsent(completeName(name), x -> new SimpleGauge(gauge)); + } + + @Override + public void unregisterGauge(String name, Gauge gauge) { + // no-op + } + + @Override + public void removeScope(String name, StatsLogger statsLogger) { + // no-op + } + + @Override + public StatsLogger scope(String name) { + return new PrometheusStatsLogger(provider, completeName(name)); + } + + private String completeName(String name) { + String completeName = scope.isEmpty() ? name : Joiner.on('_').join(scope, name); + return Collector.sanitizeMetricName(completeName); + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java new file mode 100644 index 0000000000000..6b9e662d5ff47 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java @@ -0,0 +1,152 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.stats.prometheus; + +import io.prometheus.client.Collector; +import io.prometheus.client.Collector.MetricFamilySamples; +import io.prometheus.client.Collector.MetricFamilySamples.Sample; +import io.prometheus.client.CollectorRegistry; + +import java.io.IOException; +import java.io.Writer; +import java.util.Enumeration; + +import org.apache.bookkeeper.stats.Counter; + +/** + * Logic to write metrics in Prometheus text format. + */ +public class PrometheusTextFormatUtil { + static void writeGauge(Writer w, String name, SimpleGauge gauge) { + // Example: + // # TYPE bookie_storage_entries_count gauge + // bookie_storage_entries_count 519 + try { + w.append("# TYPE ").append(name).append(" gauge\n"); + w.append(name).append(' ').append(gauge.getSample().toString()).append('\n'); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + static void writeCounter(Writer w, String name, Counter counter) { + // Example: + // # TYPE jvm_threads_started_total counter + // jvm_threads_started_total 59 + try { + w.append("# TYPE ").append(name).append(" counter\n"); + w.append(name).append(' ').append(counter.get().toString()).append('\n'); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + static void writeOpStat(Writer w, String name, DataSketchesOpStatsLogger opStat) { + // Example: + // # TYPE bookie_journal_JOURNAL_ADD_ENTRY summary + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.5",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.75",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.95",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.99",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.999",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.9999",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="1.0",} NaN + // bookie_journal_JOURNAL_ADD_ENTRY_count{success="false",} 0.0 + // bookie_journal_JOURNAL_ADD_ENTRY_sum{success="false",} 0.0 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.5",} 1.706 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.75",} 1.89 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.95",} 2.121 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.99",} 10.708 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.999",} 10.902 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.9999",} 10.902 + // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="1.0",} 10.902 + // bookie_journal_JOURNAL_ADD_ENTRY_count{success="true",} 658.0 + // bookie_journal_JOURNAL_ADD_ENTRY_sum{success="true",} 1265.0800000000002 + try { + w.append("# TYPE ").append(name).append(" summary\n"); + writeQuantile(w, opStat, name, false, 0.5); + writeQuantile(w, opStat, name, false, 0.75); + writeQuantile(w, opStat, name, false, 0.95); + writeQuantile(w, opStat, name, false, 0.99); + writeQuantile(w, opStat, name, false, 0.999); + writeQuantile(w, opStat, name, false, 0.9999); + writeQuantile(w, opStat, name, false, 1.0); + writeCount(w, opStat, name, false); + writeSum(w, opStat, name, false); + + writeQuantile(w, opStat, name, true, 0.5); + writeQuantile(w, opStat, name, true, 0.75); + writeQuantile(w, opStat, name, true, 0.95); + writeQuantile(w, opStat, name, true, 0.99); + writeQuantile(w, opStat, name, true, 0.999); + writeQuantile(w, opStat, name, true, 0.9999); + writeQuantile(w, opStat, name, true, 1.0); + writeCount(w, opStat, name, true); + writeSum(w, opStat, name, true); + + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + private static void writeQuantile(Writer w, DataSketchesOpStatsLogger opStat, String name, Boolean success, + double quantile) throws IOException { + w.append(name).append("{success=\"").append(success.toString()).append("\",quantile=\"") + .append(Double.toString(quantile)).append("\"} ") + .append(Double.toString(opStat.getQuantileValue(success, quantile))).append('\n'); + } + + private static void writeCount(Writer w, DataSketchesOpStatsLogger opStat, String name, Boolean success) + throws IOException { + w.append(name).append("_count{success=\"").append(success.toString()).append("\"} ") + .append(Long.toString(opStat.getCount(success))).append('\n'); + } + + private static void writeSum(Writer w, DataSketchesOpStatsLogger opStat, String name, Boolean success) + throws IOException { + w.append(name).append("_sum{success=\"").append(success.toString()).append("\"} ") + .append(Double.toString(opStat.getSum(success))).append('\n'); + } + + static void writeMetricsCollectedByPrometheusClient(Writer w, CollectorRegistry registry) throws IOException { + Enumeration metricFamilySamples = registry.metricFamilySamples(); + while (metricFamilySamples.hasMoreElements()) { + MetricFamilySamples metricFamily = metricFamilySamples.nextElement(); + + for (int i = 0; i < metricFamily.samples.size(); i++) { + Sample sample = metricFamily.samples.get(i); + w.write(sample.name); + w.write('{'); + for (int j = 0; j < sample.labelNames.size(); j++) { + if (j != 0) { + w.write(", "); + } + w.write(sample.labelNames.get(j)); + w.write("=\""); + w.write(sample.labelValues.get(j)); + w.write('"'); + } + + w.write("} "); + w.write(Collector.doubleToGoString(sample.value)); + w.write('\n'); + } + } + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java new file mode 100644 index 0000000000000..f6065611e3cf2 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java @@ -0,0 +1,37 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.stats.prometheus; + +import org.apache.bookkeeper.stats.Gauge; + +/** + * A {@link Gauge} implementation that forwards on the value supplier. + */ +public class SimpleGauge { + + private final Gauge gauge; + + public SimpleGauge(final Gauge gauge) { + this.gauge = gauge; + } + + Number getSample() { + return gauge.getSample(); + } +} diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 030e83893ef9e..6e33cb7be9661 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -1133,6 +1133,19 @@ public class ServiceConfiguration implements PulsarConfiguration { doc = "Add entry timeout when broker tries to publish message to bookkeeper.(0 to disable it)") private long managedLedgerAddEntryTimeoutSeconds = 0; + @FieldContext( + category = CATEGORY_STORAGE_ML, + doc = "Managed ledger prometheus stats latency rollover seconds" + ) + private int managedLedgerPrometheusStatsLatencyRolloverSeconds = 60; + + @FieldContext( + dynamic = true, + category = CATEGORY_STORAGE_ML, + doc = "Whether trace managed ledger task execution time" + ) + private boolean managedLedgerTraceTaskExecution = false; + /*** --- Load balancer --- ****/ @FieldContext( category = CATEGORY_LOAD_BALANCER, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index 09fea0c0ffc99..e1af440663a30 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -56,6 +56,8 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, managedLedgerFactoryConfig.setCacheEvictionFrequency(conf.getManagedLedgerCacheEvictionFrequency()); managedLedgerFactoryConfig.setCacheEvictionTimeThresholdMillis(conf.getManagedLedgerCacheEvictionTimeThresholdMillis()); managedLedgerFactoryConfig.setCopyEntriesInCache(conf.isManagedLedgerCacheCopyEntries()); + managedLedgerFactoryConfig.setPrometheusStatsLatencyRolloverSeconds(conf.getManagedLedgerPrometheusStatsLatencyRolloverSeconds()); + managedLedgerFactoryConfig.setTraceTaskExecution(conf.isManagedLedgerTraceTaskExecution()); this.defaultBkClient = bookkeeperProvider.create(conf, zkClient, Optional.empty(), null); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java index 5796c8e34f82c..1e0eb7b2f5a39 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java @@ -20,12 +20,18 @@ import java.io.IOException; import java.io.OutputStream; +<<<<<<< HEAD import java.util.Collection; import java.util.HashSet; import java.util.Map; import java.util.Set; +======= +import java.io.StringWriter; +import java.io.Writer; +>>>>>>> expose managed ledger bookie client metric to prometheus import java.util.Enumeration; +import org.apache.bookkeeper.stats.StatsProvider; import org.apache.pulsar.broker.PulsarService; import static org.apache.pulsar.common.stats.JvmMetrics.getJvmDirectMemoryUsed; @@ -85,6 +91,8 @@ public static void generate(PulsarService pulsar, boolean includeTopicMetrics, b generateBrokerBasicMetrics(pulsar, stream); + generateManagedLedgerBookieClientMetrics(pulsar, stream); + out.write(buf.array(), buf.arrayOffset(), buf.readableBytes()); } finally { buf.release(); @@ -155,6 +163,17 @@ private static void parseMetricsToPrometheusMetrics(Collection metrics, } } + private static void generateManagedLedgerBookieClientMetrics(PulsarService pulsar, SimpleTextOutputStream stream) { + StatsProvider statsProvider = pulsar.getManagedLedgerClientFactory().getManagedLedgerFactory().getStatsProvider(); + try { + Writer writer = new StringWriter(); + statsProvider.writeAllMetrics(writer); + stream.write(writer.toString()); + } catch (IOException e) { + // nop + } + } + private static void generateSystemMetrics(SimpleTextOutputStream stream, String cluster) { Enumeration metricFamilySamples = CollectorRegistry.defaultRegistry.metricFamilySamples(); while (metricFamilySamples.hasMoreElements()) { From 101340445049fe3515a793e3c52ee9c27a53af68 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sat, 25 Apr 2020 11:19:45 +0800 Subject: [PATCH 02/12] format metric name --- .../broker/stats/prometheus/PrometheusMetricsGenerator.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java index 1e0eb7b2f5a39..15e28125a8cea 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java @@ -20,15 +20,12 @@ import java.io.IOException; import java.io.OutputStream; -<<<<<<< HEAD import java.util.Collection; import java.util.HashSet; import java.util.Map; import java.util.Set; -======= import java.io.StringWriter; import java.io.Writer; ->>>>>>> expose managed ledger bookie client metric to prometheus import java.util.Enumeration; import org.apache.bookkeeper.stats.StatsProvider; From 703ccfe3e755e365643109a36db0d858ce7c71f0 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sat, 25 Apr 2020 11:22:29 +0800 Subject: [PATCH 03/12] format name --- .../apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index d8e405e09bf47..071432cbc1e18 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -58,7 +58,7 @@ public class ManagedLedgerFactoryConfig { private boolean traceTaskExecution = false; /** - * managed ledger prometheus stats Latency Rollover Seconds + * Managed ledger prometheus stats Latency Rollover Seconds */ private int prometheusStatsLatencyRolloverSeconds = 60; } From 881d71fb265e891df31484c27b391e0fd327963e Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sat, 25 Apr 2020 17:34:08 +0800 Subject: [PATCH 04/12] format code --- .../impl/ManagedLedgerFactoryImpl.java | 4 +- .../prometheus/PrometheusMetricsProvider.java | 10 ++--- .../prometheus/PrometheusTextFormatUtil.java | 42 +++++++++---------- 3 files changed, 29 insertions(+), 27 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index a08366e0cb55d..11cc10c101be6 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -158,7 +158,9 @@ public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolic private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, boolean isBookkeeperManaged, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config) throws Exception { Configuration configuration = new ClientConfiguration(); - configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS, config.getPrometheusStatsLatencyRolloverSeconds()); + + configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS + , config.getPrometheusStatsLatencyRolloverSeconds()); statsProvider.start(configuration); statsLogger = statsProvider.getStatsLogger("pulsar_bookie_client"); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java index b9c53bd9b37c6..0f1dabbdfbfae 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java @@ -21,15 +21,16 @@ import com.google.common.annotations.VisibleForTesting; import io.netty.util.concurrent.DefaultThreadFactory; -import io.netty.util.internal.PlatformDependent; import io.prometheus.client.Collector; import io.prometheus.client.CollectorRegistry; import java.io.IOException; import java.io.Writer; -import java.lang.reflect.Field; -import java.util.concurrent.*; -import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.ConcurrentSkipListMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import org.apache.bookkeeper.stats.CachingStatsProvider; import org.apache.bookkeeper.stats.StatsLogger; @@ -149,6 +150,5 @@ private void registerMetrics(Collector collector) { } } - private static final Logger log = LoggerFactory.getLogger(org.apache.bookkeeper.stats.prometheus.PrometheusMetricsProvider.class); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java index 6b9e662d5ff47..f180e9db10a83 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java @@ -35,8 +35,8 @@ public class PrometheusTextFormatUtil { static void writeGauge(Writer w, String name, SimpleGauge gauge) { // Example: - // # TYPE bookie_storage_entries_count gauge - // bookie_storage_entries_count 519 + // # TYPE bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 gauge + // pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 1044057 try { w.append("# TYPE ").append(name).append(" gauge\n"); w.append(name).append(' ').append(gauge.getSample().toString()).append('\n'); @@ -59,25 +59,25 @@ static void writeCounter(Writer w, String name, Counter counter) { static void writeOpStat(Writer w, String name, DataSketchesOpStatsLogger opStat) { // Example: - // # TYPE bookie_journal_JOURNAL_ADD_ENTRY summary - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.5",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.75",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.95",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.99",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.999",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="0.9999",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY{success="false",quantile="1.0",} NaN - // bookie_journal_JOURNAL_ADD_ENTRY_count{success="false",} 0.0 - // bookie_journal_JOURNAL_ADD_ENTRY_sum{success="false",} 0.0 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.5",} 1.706 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.75",} 1.89 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.95",} 2.121 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.99",} 10.708 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.999",} 10.902 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="0.9999",} 10.902 - // bookie_journal_JOURNAL_ADD_ENTRY{success="true",quantile="1.0",} 10.902 - // bookie_journal_JOURNAL_ADD_ENTRY_count{success="true",} 658.0 - // bookie_journal_JOURNAL_ADD_ENTRY_sum{success="true",} 1265.0800000000002 + // # TYPE pulsar_bookie_client_bookkeeper_ml_workers_task_queued summary + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.5"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.75"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.95"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.99"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.999"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.9999"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="1.0"} -Infinity + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_count{success="false"} 0 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{success="false"} 0.0 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="0.5"} 0.031 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="0.75"} 0.043 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="0.95"} 0.061 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="0.99"} 0.064 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="0.999"} 0.073 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="0.9999"} 0.073 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="true",quantile="1.0"} 0.552 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_count{success="true"} 40911432 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{success="true"} 527.0 try { w.append("# TYPE ").append(name).append(" summary\n"); writeQuantile(w, opStat, name, false, 0.5); From 6d7075b70c6c966057a4d2170c0c5dd51316f147 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sat, 25 Apr 2020 18:29:31 +0800 Subject: [PATCH 05/12] update document --- .../mledger/impl/ManagedLedgerFactoryImpl.java | 1 - site2/docs/reference-metrics.md | 17 +++++++++++++++++ 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 11cc10c101be6..232f39628c625 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -158,7 +158,6 @@ public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolic private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, boolean isBookkeeperManaged, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config) throws Exception { Configuration configuration = new ClientConfiguration(); - configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS , config.getPrometheusStatsLatencyRolloverSeconds()); statsProvider.start(configuration); diff --git a/site2/docs/reference-metrics.md b/site2/docs/reference-metrics.md index 11cdc21fed003..9c0a2633a6700 100644 --- a/site2/docs/reference-metrics.md +++ b/site2/docs/reference-metrics.md @@ -107,6 +107,7 @@ Broker has the following kinds of metrics: * [BundleSplit metrics](#bundlesplit-metrics) * [Subscription metrics](#subscription-metrics) * [Consumer metrics](#consumer-metrics) +* [ManagedLedger bookie client metrics](#managed-ledger-bookie-client-metrics) ### Namespace metrics @@ -319,6 +320,22 @@ All the consumer metrics are labelled with the following labels: | pulsar_consumer_msg_throughput_out | Gauge | The total message dispatch throughput for a consumer (bytes/second). | | pulsar_consumer_available_permits | Gauge | The available permits for for a consumer. | +### Managed ledger bookie client metrics + +All the managed ledger bookie client metrics list as follows: + +| Name | Type | Description | +| --- | --- | --- | +| pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_* | Gauge | The number of tasks the scheduler executor execute completed.
The number of metrics determined by the scheduler executor thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| +| pulsar_bookie_client_bookkeeper_ml_scheduler_queue_* | Gauge | The number of tasks queued in the scheduler executor's queue.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| +| pulsar_bookie_client_bookkeeper_ml_scheduler_total_tasks_* | Gauge | The total number of tasks the scheduler executor received.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| +| pulsar_bookie_client_bookkeeper_ml_workers_completed_tasks_* | Gauge | The number of tasks the worker executor execute completed.
The number of metrics determined by the number of worker task thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`
| +| pulsar_bookie_client_bookkeeper_ml_workers_queue_* | Gauge | The number of tasks queued in the worker executor's queue.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`.
| +| pulsar_bookie_client_bookkeeper_ml_workers_total_tasks_* | Gauge | The total number of tasks the worker executor received.
The number of metrics determined by worker executor's thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`.
| +| pulsar_bookie_client_bookkeeper_ml_scheduler_task_execution | Summary | The scheduler task execution latency calculated in milliseconds. | +| pulsar_bookie_client_bookkeeper_ml_scheduler_task_queued | Summary | The scheduler task queued latency calculated in milliseconds. | +| pulsar_bookie_client_bookkeeper_ml_workers_task_execution | Summary | The worker task execution latency calculated in milliseconds. | +| pulsar_bookie_client_bookkeeper_ml_workers_task_queued | Summary | The worker task queued latency calculated in milliseconds. | ## Monitor You can [set up a Prometheus instance](https://prometheus.io/) to collect all the metrics exposed at Pulsar components and set up From 6193c5a286750d66e856676fe77391de9280be9d Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Mon, 27 Apr 2020 13:30:08 +0800 Subject: [PATCH 06/12] fix a bug and turn on traceTaskExecution --- conf/broker.conf | 2 +- conf/standalone.conf | 2 +- .../mledger/ManagedLedgerFactoryConfig.java | 2 +- .../prometheus/PrometheusMetricsProvider.java | 23 ------------------- .../pulsar/broker/ServiceConfiguration.java | 2 +- 5 files changed, 4 insertions(+), 27 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index bc49086acb8d5..cca5f48f9df19 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -684,7 +684,7 @@ managedLedgerAddEntryTimeoutSeconds=0 managedLedgerPrometheusStatsLatencyRolloverSeconds=60 # Whether trace managed ledger task execution time -managedLedgerTraceTaskExecution=false +managedLedgerTraceTaskExecution=true ### --- Load balancer --- ### diff --git a/conf/standalone.conf b/conf/standalone.conf index 28cf754bf81df..674b71f95f752 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -453,7 +453,7 @@ managedLedgerUnackedRangesOpenCacheSetEnabled=true managedLedgerPrometheusStatsLatencyRolloverSeconds=60 # Whether trace managed ledger task execution time -managedLedgerTraceTaskExecution=false +managedLedgerTraceTaskExecution=true ### --- Load balancer --- ### diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index 071432cbc1e18..469aff4783610 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -55,7 +55,7 @@ public class ManagedLedgerFactoryConfig { /** * Whether trace managed ledger task execution time */ - private boolean traceTaskExecution = false; + private boolean traceTaskExecution = true; /** * Managed ledger prometheus stats Latency Rollover Seconds diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java index 0f1dabbdfbfae..c5258489b26a3 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java @@ -22,7 +22,6 @@ import io.netty.util.concurrent.DefaultThreadFactory; import io.prometheus.client.Collector; -import io.prometheus.client.CollectorRegistry; import java.io.IOException; import java.io.Writer; @@ -37,8 +36,6 @@ import org.apache.bookkeeper.stats.StatsProvider; import org.apache.commons.configuration.Configuration; import org.apache.commons.lang.StringUtils; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * A Prometheus based {@link StatsProvider} implementation. @@ -49,7 +46,6 @@ public class PrometheusMetricsProvider implements StatsProvider { public static final String PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS = "prometheusStatsLatencyRolloverSeconds"; public static final int DEFAULT_PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS = 60; - final CollectorRegistry registry; private final CachingStatsProvider cachingStatsProvider; /** @@ -60,11 +56,6 @@ public class PrometheusMetricsProvider implements StatsProvider { final ConcurrentMap opStats = new ConcurrentSkipListMap<>(); public PrometheusMetricsProvider() { - this(CollectorRegistry.defaultRegistry); - } - - public PrometheusMetricsProvider(CollectorRegistry registry) { - this.registry = registry; this.cachingStatsProvider = new CachingStatsProvider(new StatsProvider() { @Override public void start(Configuration conf) { @@ -120,8 +111,6 @@ public StatsLogger getStatsLogger(String scope) { @Override public void writeAllMetrics(Writer writer) throws IOException { - PrometheusTextFormatUtil.writeMetricsCollectedByPrometheusClient(writer, registry); - gauges.forEach((name, gauge) -> PrometheusTextFormatUtil.writeGauge(writer, name, gauge)); counters.forEach((name, counter) -> PrometheusTextFormatUtil.writeCounter(writer, name, counter)); opStats.forEach((name, opStatLogger) -> PrometheusTextFormatUtil.writeOpStat(writer, name, opStatLogger)); @@ -139,16 +128,4 @@ void rotateLatencyCollection() { }); } - private void registerMetrics(Collector collector) { - try { - collector.register(registry); - } catch (Exception e) { - // Ignore if these were already registered - if (log.isDebugEnabled()) { - log.debug("Failed to register Prometheus collector exports", e); - } - } - } - - private static final Logger log = LoggerFactory.getLogger(org.apache.bookkeeper.stats.prometheus.PrometheusMetricsProvider.class); } diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 6e33cb7be9661..9d26553a85771 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -1144,7 +1144,7 @@ public class ServiceConfiguration implements PulsarConfiguration { category = CATEGORY_STORAGE_ML, doc = "Whether trace managed ledger task execution time" ) - private boolean managedLedgerTraceTaskExecution = false; + private boolean managedLedgerTraceTaskExecution = true; /*** --- Load balancer --- ****/ @FieldContext( From cde977f18890af1bda89ff5feccd0b197c0c3c9e Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Fri, 1 May 2020 23:18:44 +0800 Subject: [PATCH 07/12] update metric labels --- .../mledger/ManagedLedgerFactoryConfig.java | 5 ++ .../impl/ManagedLedgerFactoryImpl.java | 1 + .../prometheus/PrometheusMetricsProvider.java | 10 ++- .../prometheus/PrometheusTextFormatUtil.java | 69 ++++++++++--------- .../broker/ManagedLedgerClientFactory.java | 1 + .../broker/stats/PrometheusMetricsTest.java | 44 ++++++++++++ site2/docs/reference-metrics.md | 5 +- 7 files changed, 99 insertions(+), 36 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java index 469aff4783610..e42befc33a410 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactoryConfig.java @@ -61,4 +61,9 @@ public class ManagedLedgerFactoryConfig { * Managed ledger prometheus stats Latency Rollover Seconds */ private int prometheusStatsLatencyRolloverSeconds = 60; + + /** + * cluster name for prometheus stats + */ + private String clusterName; } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 232f39628c625..1054a97bf8e4c 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -160,6 +160,7 @@ private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPoli Configuration configuration = new ClientConfiguration(); configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS , config.getPrometheusStatsLatencyRolloverSeconds()); + configuration.addProperty(PrometheusMetricsProvider.CLUSTER_NAME, config.getClusterName()); statsProvider.start(configuration); statsLogger = statsProvider.getStatsLogger("pulsar_bookie_client"); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java index c5258489b26a3..fc49289063ffa 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java @@ -45,7 +45,10 @@ public class PrometheusMetricsProvider implements StatsProvider { public static final String PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS = "prometheusStatsLatencyRolloverSeconds"; public static final int DEFAULT_PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS = 60; + public static final String CLUSTER_NAME = "cluster"; + public static final String DEFAULT_CLUSTER_NAME = "pulsar"; + private String cluster; private final CachingStatsProvider cachingStatsProvider; /** @@ -93,6 +96,7 @@ public void start(Configuration conf) { int latencyRolloverSeconds = conf.getInt(PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS, DEFAULT_PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS); + cluster = conf.getString(CLUSTER_NAME, DEFAULT_CLUSTER_NAME); executor.scheduleAtFixedRate(() -> { rotateLatencyCollection(); @@ -111,9 +115,9 @@ public StatsLogger getStatsLogger(String scope) { @Override public void writeAllMetrics(Writer writer) throws IOException { - gauges.forEach((name, gauge) -> PrometheusTextFormatUtil.writeGauge(writer, name, gauge)); - counters.forEach((name, counter) -> PrometheusTextFormatUtil.writeCounter(writer, name, counter)); - opStats.forEach((name, opStatLogger) -> PrometheusTextFormatUtil.writeOpStat(writer, name, opStatLogger)); + gauges.forEach((name, gauge) -> PrometheusTextFormatUtil.writeGauge(writer, name, cluster, gauge)); + counters.forEach((name, counter) -> PrometheusTextFormatUtil.writeCounter(writer, name, cluster, counter)); + opStats.forEach((name, opStatLogger) -> PrometheusTextFormatUtil.writeOpStat(writer, name, cluster, opStatLogger)); } @Override diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java index f180e9db10a83..af9d378b64642 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java @@ -33,31 +33,33 @@ * Logic to write metrics in Prometheus text format. */ public class PrometheusTextFormatUtil { - static void writeGauge(Writer w, String name, SimpleGauge gauge) { + static void writeGauge(Writer w, String name, String cluster, SimpleGauge gauge) { // Example: // # TYPE bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 gauge // pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 1044057 try { w.append("# TYPE ").append(name).append(" gauge\n"); - w.append(name).append(' ').append(gauge.getSample().toString()).append('\n'); + w.append(name).append("{cluster=\"").append(cluster).append("\"}") + .append(' ').append(gauge.getSample().toString()).append('\n'); } catch (IOException e) { throw new RuntimeException(e); } } - static void writeCounter(Writer w, String name, Counter counter) { + static void writeCounter(Writer w, String name, String cluster, Counter counter) { // Example: // # TYPE jvm_threads_started_total counter // jvm_threads_started_total 59 try { w.append("# TYPE ").append(name).append(" counter\n"); - w.append(name).append(' ').append(counter.get().toString()).append('\n'); + w.append(name).append("{cluster=\"").append(cluster).append("\"}") + .append(' ').append(counter.get().toString()).append('\n'); } catch (IOException e) { throw new RuntimeException(e); } } - static void writeOpStat(Writer w, String name, DataSketchesOpStatsLogger opStat) { + static void writeOpStat(Writer w, String name, String cluster, DataSketchesOpStatsLogger opStat) { // Example: // # TYPE pulsar_bookie_client_bookkeeper_ml_workers_task_queued summary // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{success="false",quantile="0.5"} NaN @@ -80,47 +82,50 @@ static void writeOpStat(Writer w, String name, DataSketchesOpStatsLogger opStat) // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{success="true"} 527.0 try { w.append("# TYPE ").append(name).append(" summary\n"); - writeQuantile(w, opStat, name, false, 0.5); - writeQuantile(w, opStat, name, false, 0.75); - writeQuantile(w, opStat, name, false, 0.95); - writeQuantile(w, opStat, name, false, 0.99); - writeQuantile(w, opStat, name, false, 0.999); - writeQuantile(w, opStat, name, false, 0.9999); - writeQuantile(w, opStat, name, false, 1.0); - writeCount(w, opStat, name, false); - writeSum(w, opStat, name, false); + writeQuantile(w, opStat, name, cluster,false, 0.5); + writeQuantile(w, opStat, name, cluster, false, 0.75); + writeQuantile(w, opStat, name, cluster,false, 0.95); + writeQuantile(w, opStat, name, cluster,false, 0.99); + writeQuantile(w, opStat, name, cluster,false, 0.999); + writeQuantile(w, opStat, name, cluster,false, 0.9999); + writeQuantile(w, opStat, name, cluster,false, 1.0); + writeCount(w, opStat, name, cluster, false); + writeSum(w, opStat, name, cluster, false); - writeQuantile(w, opStat, name, true, 0.5); - writeQuantile(w, opStat, name, true, 0.75); - writeQuantile(w, opStat, name, true, 0.95); - writeQuantile(w, opStat, name, true, 0.99); - writeQuantile(w, opStat, name, true, 0.999); - writeQuantile(w, opStat, name, true, 0.9999); - writeQuantile(w, opStat, name, true, 1.0); - writeCount(w, opStat, name, true); - writeSum(w, opStat, name, true); + writeQuantile(w, opStat, name, cluster,true, 0.5); + writeQuantile(w, opStat, name, cluster,true, 0.75); + writeQuantile(w, opStat, name, cluster,true, 0.95); + writeQuantile(w, opStat, name, cluster,true, 0.99); + writeQuantile(w, opStat, name, cluster,true, 0.999); + writeQuantile(w, opStat, name, cluster,true, 0.9999); + writeQuantile(w, opStat, name, cluster,true, 1.0); + writeCount(w, opStat, name, cluster, true); + writeSum(w, opStat, name, cluster, true); } catch (IOException e) { throw new RuntimeException(e); } } - private static void writeQuantile(Writer w, DataSketchesOpStatsLogger opStat, String name, Boolean success, - double quantile) throws IOException { - w.append(name).append("{success=\"").append(success.toString()).append("\",quantile=\"") + private static void writeQuantile(Writer w, DataSketchesOpStatsLogger opStat, String name, String cluster, + Boolean success, double quantile) throws IOException { + w.append(name).append("{cluster=\"").append(cluster).append("\", success=\"") + .append(success.toString()).append("\",quantile=\"") .append(Double.toString(quantile)).append("\"} ") .append(Double.toString(opStat.getQuantileValue(success, quantile))).append('\n'); } - private static void writeCount(Writer w, DataSketchesOpStatsLogger opStat, String name, Boolean success) - throws IOException { - w.append(name).append("_count{success=\"").append(success.toString()).append("\"} ") + private static void writeCount(Writer w, DataSketchesOpStatsLogger opStat, String name, String cluster, + Boolean success) throws IOException { + w.append(name).append("_count{cluster=\"").append(cluster).append("\", success=\"") + .append(success.toString()).append("\"} ") .append(Long.toString(opStat.getCount(success))).append('\n'); } - private static void writeSum(Writer w, DataSketchesOpStatsLogger opStat, String name, Boolean success) - throws IOException { - w.append(name).append("_sum{success=\"").append(success.toString()).append("\"} ") + private static void writeSum(Writer w, DataSketchesOpStatsLogger opStat, String name, String cluster, + Boolean success) throws IOException { + w.append(name).append("_sum{cluster=\"").append(cluster).append("\", success=\"") + .append(success.toString()).append("\"} ") .append(Double.toString(opStat.getSum(success))).append('\n'); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index e1af440663a30..5a3a9cabe90b8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -58,6 +58,7 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, managedLedgerFactoryConfig.setCopyEntriesInCache(conf.isManagedLedgerCacheCopyEntries()); managedLedgerFactoryConfig.setPrometheusStatsLatencyRolloverSeconds(conf.getManagedLedgerPrometheusStatsLatencyRolloverSeconds()); managedLedgerFactoryConfig.setTraceTaskExecution(conf.isManagedLedgerTraceTaskExecution()); + managedLedgerFactoryConfig.setClusterName(conf.getClusterName()); this.defaultBkClient = bookkeeperProvider.create(conf, zkClient, Optional.empty(), null); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java index b3246c12ff043..a1472e7ba4833 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java @@ -309,6 +309,50 @@ public void testManagedLedgerStats() throws Exception { p2.close(); } + @Test + public void testManagedLedgerBookieClientStats() throws Exception { + Producer p1 = pulsarClient.newProducer().topic("persistent://my-property/use/my-ns/my-topic1").create(); + Producer p2 = pulsarClient.newProducer().topic("persistent://my-property/use/my-ns/my-topic2").create(); + for (int i = 0; i < 10; i++) { + String message = "my-message-" + i; + p1.send(message.getBytes()); + p2.send(message.getBytes()); + } + + ByteArrayOutputStream statsOut = new ByteArrayOutputStream(); + PrometheusMetricsGenerator.generate(pulsar, false, false, statsOut); + String metricsStr = new String(statsOut.toByteArray()); + + Multimap metrics = parseMetrics(metricsStr); + + metrics.entries().forEach(e -> + System.out.println(e.getKey() + ": " + e.getValue()) + ); + + List cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0"); + assertEquals(cm.size(), 1); + assertEquals(cm.get(0).tags.get("cluster"), "test"); + + cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_scheduler_queue_0"); + assertEquals(cm.size(), 1); + assertEquals(cm.get(0).tags.get("cluster"), "test"); + + cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_scheduler_total_tasks_0"); + assertEquals(cm.size(), 1); + assertEquals(cm.get(0).tags.get("cluster"), "test"); + + cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_workers_completed_tasks_0"); + assertEquals(cm.size(), 1); + assertEquals(cm.get(0).tags.get("cluster"), "test"); + + cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_workers_task_execution_count"); + assertEquals(cm.size(), 2); + assertEquals(cm.get(0).tags.get("cluster"), "test"); + + p1.close(); + p2.close(); + } + /** * Hacky parsing of Prometheus text format. Sould be good enough for unit tests */ diff --git a/site2/docs/reference-metrics.md b/site2/docs/reference-metrics.md index 9c0a2633a6700..03a4bb757eadb 100644 --- a/site2/docs/reference-metrics.md +++ b/site2/docs/reference-metrics.md @@ -322,7 +322,9 @@ All the consumer metrics are labelled with the following labels: ### Managed ledger bookie client metrics -All the managed ledger bookie client metrics list as follows: +All the managed ledger bookie client metrics labelled with the following labels: + +- *cluster*: `cluster=${pulsar_cluster}`. `${pulsar_cluster}` is the cluster name that you configured in `broker.conf`. | Name | Type | Description | | --- | --- | --- | @@ -336,6 +338,7 @@ All the managed ledger bookie client metrics list as follows: | pulsar_bookie_client_bookkeeper_ml_scheduler_task_queued | Summary | The scheduler task queued latency calculated in milliseconds. | | pulsar_bookie_client_bookkeeper_ml_workers_task_execution | Summary | The worker task execution latency calculated in milliseconds. | | pulsar_bookie_client_bookkeeper_ml_workers_task_queued | Summary | The worker task queued latency calculated in milliseconds. | + ## Monitor You can [set up a Prometheus instance](https://prometheus.io/) to collect all the metrics exposed at Pulsar components and set up From dba0bc12ac0e5566d232b76323baf535e60a4a17 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sat, 2 May 2020 09:54:25 +0800 Subject: [PATCH 08/12] simplify code --- .../stats/prometheus/LongAdderCounter.java | 58 ------------------- .../prometheus/PrometheusMetricsProvider.java | 1 + .../prometheus/PrometheusStatsLogger.java | 2 + .../prometheus/PrometheusTextFormatUtil.java | 40 ++++++------- 4 files changed, 23 insertions(+), 78 deletions(-) delete mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java deleted file mode 100644 index 44fc4fafb6e2e..0000000000000 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/LongAdderCounter.java +++ /dev/null @@ -1,58 +0,0 @@ -/** - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - */ -package org.apache.bookkeeper.mledger.stats.prometheus; - -import java.util.concurrent.atomic.LongAdder; - -import org.apache.bookkeeper.stats.Counter; - -/** - * {@link Counter} implementation based on {@link LongAdder}. - * - *

LongAdder keeps a counter per-thread and then aggregates to get the result, in order to avoid contention between - * multiple threads. - */ -public class LongAdderCounter implements Counter { - private final LongAdder counter = new LongAdder(); - - @Override - public void clear() { - counter.reset(); - } - - @Override - public void inc() { - counter.increment(); - } - - @Override - public void dec() { - counter.decrement(); - } - - @Override - public void add(long delta) { - counter.add(delta); - } - - @Override - public Long get() { - return counter.sum(); - } -} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java index fc49289063ffa..dc8b980b116c0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java @@ -34,6 +34,7 @@ import org.apache.bookkeeper.stats.CachingStatsProvider; import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.bookkeeper.stats.prometheus.LongAdderCounter; import org.apache.commons.configuration.Configuration; import org.apache.commons.lang.StringUtils; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java index 295d4575ca0ca..ca8180d0b67b0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java @@ -27,6 +27,8 @@ import org.apache.bookkeeper.stats.OpStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; +import org.apache.bookkeeper.stats.prometheus.LongAdderCounter; + /** * A {@code Prometheus} based {@link StatsLogger} implementation. */ diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java index af9d378b64642..30f45a601a134 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java @@ -36,7 +36,7 @@ public class PrometheusTextFormatUtil { static void writeGauge(Writer w, String name, String cluster, SimpleGauge gauge) { // Example: // # TYPE bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 gauge - // pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 1044057 + // pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0{cluster="test"} 1044057 try { w.append("# TYPE ").append(name).append(" gauge\n"); w.append(name).append("{cluster=\"").append(cluster).append("\"}") @@ -49,7 +49,7 @@ static void writeGauge(Writer w, String name, String cluster, SimpleGauge Date: Mon, 4 May 2020 19:29:45 +0800 Subject: [PATCH 09/12] change prometheus stats logger to ManagedLedgerFactory and support configuration for stats expose --- conf/broker.conf | 3 ++ conf/standalone.conf | 3 ++ .../mledger/ManagedLedgerFactory.java | 6 --- .../impl/ManagedLedgerFactoryImpl.java | 39 +++++++---------- .../pulsar/broker/ServiceConfiguration.java | 6 +++ .../broker/BookKeeperClientFactory.java | 5 +++ .../broker/BookKeeperClientFactoryImpl.java | 13 +++++- .../broker/ManagedLedgerClientFactory.java | 27 +++++++++++- .../PrometheusMetricsGenerator.java | 7 ++- .../metrics}/DataSketchesOpStatsLogger.java | 8 ++-- .../metrics}/PrometheusMetricsProvider.java | 17 +++----- .../metrics}/PrometheusStatsLogger.java | 5 +-- .../metrics}/PrometheusTextFormatUtil.java | 43 +++++++++---------- .../prometheus/metrics}/SimpleGauge.java | 2 +- .../broker/MockedBookKeeperClientFactory.java | 8 ++++ .../auth/MockedPulsarServiceBaseTest.java | 9 ++++ site2/docs/reference-metrics.md | 20 ++++----- 17 files changed, 134 insertions(+), 87 deletions(-) rename {managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus => pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics}/DataSketchesOpStatsLogger.java (99%) rename {managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus => pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics}/PrometheusMetricsProvider.java (94%) rename {managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus => pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics}/PrometheusStatsLogger.java (97%) rename {managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus => pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics}/PrometheusTextFormatUtil.java (87%) rename {managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus => pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics}/SimpleGauge.java (95%) diff --git a/conf/broker.conf b/conf/broker.conf index cca5f48f9df19..3a72ba53fcb88 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -581,6 +581,9 @@ bookkeeperDiskWeightBasedPlacementEnabled=false # A value of '0' disables sending any explicit LACs. Default is 0. bookkeeperExplicitLacIntervalInMills=0 +# Expose bookkeeper client managed ledger stats to prometheus. default is false +# bookkeeperClientExposeStatsToPrometheus=false + ### --- Managed Ledger --- ### # Number of bookies to use when creating a ledger diff --git a/conf/standalone.conf b/conf/standalone.conf index 674b71f95f752..8c5657645e900 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -355,6 +355,9 @@ bookkeeperDiskWeightBasedPlacementEnabled=false # A value of '0' disables sending any explicit LACs. Default is 0. bookkeeperExplicitLacIntervalInMills=0 +# Expose bookkeeper client managed ledger stats to prometheus. default is false +# bookkeeperClientExposeStatsToPrometheus=false + ### --- Managed Ledger --- ### # Number of bookies to use when creating a ledger diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java index 286b5f72761b2..2f9b2f6ddbcb0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerFactory.java @@ -141,10 +141,4 @@ void asyncOpenReadOnlyCursor(String managedLedgerName, Position startPosition, M */ void shutdown() throws InterruptedException, ManagedLedgerException; - /** - * Get managed ledger stats provider. - * - * @return StatsProvider - */ - StatsProvider getStatsProvider(); } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 1054a97bf8e4c..f590e1e27bc2a 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -112,9 +112,6 @@ public class ManagedLedgerFactoryImpl implements ManagedLedgerFactory { private static final int StatsPeriodSeconds = 60; - private StatsLogger statsLogger = NullStatsLogger.INSTANCE; - private StatsProvider statsProvider = new PrometheusMetricsProvider(); - public ManagedLedgerFactoryImpl(ClientConfiguration bkClientConfiguration, String zkConnection) throws Exception { this(bkClientConfiguration, zkConnection, new ManagedLedgerFactoryConfig()); } @@ -130,7 +127,8 @@ public ManagedLedgerFactoryImpl(ClientConfiguration bkClientConfiguration, Manag private ManagedLedgerFactoryImpl(ZooKeeper zkc, ClientConfiguration bkClientConfiguration, ManagedLedgerFactoryConfig config) throws Exception { - this(new DefaultBkFactory(bkClientConfiguration, zkc), true /* isBookkeeperManaged */, zkc, config); + this(new DefaultBkFactory(bkClientConfiguration, zkc), true /* isBookkeeperManaged */, + zkc, config, NullStatsLogger.INSTANCE); } private ManagedLedgerFactoryImpl(ClientConfiguration clientConfiguration, String zkConnection, ManagedLedgerFactoryConfig config) throws Exception { @@ -138,7 +136,7 @@ private ManagedLedgerFactoryImpl(ClientConfiguration clientConfiguration, String true, ZooKeeperClient.newBuilder() .connectString(zkConnection) - .sessionTimeoutMs(clientConfiguration.getZkTimeout()).build(), config); + .sessionTimeoutMs(clientConfiguration.getZkTimeout()).build(), config, NullStatsLogger.INSTANCE); } public ManagedLedgerFactoryImpl(BookKeeper bookKeeper, ZooKeeper zooKeeper) throws Exception { @@ -147,24 +145,25 @@ public ManagedLedgerFactoryImpl(BookKeeper bookKeeper, ZooKeeper zooKeeper) thro public ManagedLedgerFactoryImpl(BookKeeper bookKeeper, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config) throws Exception { - this((policyConfig) -> bookKeeper, false /* isBookkeeperManaged */, zooKeeper, config); + this((policyConfig) -> bookKeeper, false /* isBookkeeperManaged */, + zooKeeper, config, NullStatsLogger.INSTANCE); } - public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config) + public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, + ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config) throws Exception { - this(bookKeeperGroupFactory, false /* isBookkeeperManaged */, zooKeeper, config); + this(bookKeeperGroupFactory, false /* isBookkeeperManaged */, zooKeeper, config, NullStatsLogger.INSTANCE); } - private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, boolean isBookkeeperManaged, ZooKeeper zooKeeper, - ManagedLedgerFactoryConfig config) throws Exception { - Configuration configuration = new ClientConfiguration(); - configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS - , config.getPrometheusStatsLatencyRolloverSeconds()); - configuration.addProperty(PrometheusMetricsProvider.CLUSTER_NAME, config.getClusterName()); - statsProvider.start(configuration); - - statsLogger = statsProvider.getStatsLogger("pulsar_bookie_client"); + public ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, + ZooKeeper zooKeeper, ManagedLedgerFactoryConfig config, StatsLogger statsLogger) + throws Exception { + this(bookKeeperGroupFactory, false /* isBookkeeperManaged */, zooKeeper, config, statsLogger); + } + private ManagedLedgerFactoryImpl(BookkeeperFactoryForCustomEnsemblePlacementPolicy bookKeeperGroupFactory, + boolean isBookkeeperManaged, ZooKeeper zooKeeper, + ManagedLedgerFactoryConfig config, StatsLogger statsLogger) throws Exception { scheduledExecutor = OrderedScheduler.newSchedulerBuilder() .numThreads(config.getNumManagedLedgerSchedulerThreads()) .statsLogger(statsLogger) @@ -217,11 +216,6 @@ public BookKeeper get(EnsemblePlacementPolicyConfig policy) { } } - @Override - public StatsProvider getStatsProvider() { - return statsProvider; - } - private synchronized void refreshStats() { long now = System.nanoTime(); long period = now - lastStatTimestamp; @@ -487,7 +481,6 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) { scheduledExecutor.shutdownNow(); orderedExecutor.shutdownNow(); cacheEvictionExecutor.shutdownNow(); - statsProvider.stop(); entryCacheManager.clear(); try { diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 9d26553a85771..9071e15f3347c 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -942,6 +942,12 @@ public class ServiceConfiguration implements PulsarConfiguration { @FieldContext(category = CATEGORY_STORAGE_BK, doc = "Set the interval to check the need for sending an explicit LAC") private int bookkeeperExplicitLacIntervalInMills = 0; + @FieldContext( + category = CATEGORY_STORAGE_BK, + doc = "whether expose managed ledger client stats to prometheus" + ) + private boolean bookkeeperClientExposeStatsToPrometheus = false; + /**** --- Managed Ledger --- ****/ @FieldContext( minValue = 1, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactory.java index c0c43ab8eb9e5..476857a0382fe 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactory.java @@ -24,6 +24,7 @@ import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.EnsemblePlacementPolicy; +import org.apache.bookkeeper.stats.StatsLogger; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.zookeeper.ZooKeeper; @@ -35,5 +36,9 @@ BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, Optional> ensemblePlacementPolicyClass, Map ensemblePlacementPolicyProperties) throws IOException; + BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, + Optional> ensemblePlacementPolicyClass, + Map ensemblePlacementPolicyProperties, + StatsLogger statsLogger) throws IOException; void close(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java index 095ecd03631ff..79ea3be5562d2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java @@ -37,6 +37,8 @@ import org.apache.bookkeeper.client.RackawareEnsemblePlacementPolicy; import org.apache.bookkeeper.client.RegionAwareEnsemblePlacementPolicy; import org.apache.bookkeeper.conf.ClientConfiguration; +import org.apache.bookkeeper.stats.NullStatsLogger; +import org.apache.bookkeeper.stats.StatsLogger; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.zookeeper.ZkBookieRackAffinityMapping; @@ -54,7 +56,15 @@ public class BookKeeperClientFactoryImpl implements BookKeeperClientFactory { @Override public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, - Optional> ensemblePlacementPolicyClass, Map properties) throws IOException { + Optional> ensemblePlacementPolicyClass, + Map properties) throws IOException { + return create(conf, zkClient, ensemblePlacementPolicyClass, properties, NullStatsLogger.INSTANCE); + } + + @Override + public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, + Optional> ensemblePlacementPolicyClass, + Map properties, StatsLogger statsLogger) throws IOException { ClientConfiguration bkConf = createBkClientConfiguration(conf); if (properties != null) { properties.forEach((key, value) -> bkConf.setProperty(key, value)); @@ -67,6 +77,7 @@ public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, try { return BookKeeper.forConfig(bkConf) .allocator(PulsarByteBufAllocator.DEFAULT) + .statsLogger(statsLogger) .build(); } catch (InterruptedException | BKException e) { throw new IOException(e); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java index 5a3a9cabe90b8..d96b32b25de18 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/ManagedLedgerClientFactory.java @@ -25,12 +25,18 @@ import java.util.concurrent.RejectedExecutionException; import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl.BookkeeperFactoryForCustomEnsemblePlacementPolicy; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl.EnsemblePlacementPolicyConfig; +import org.apache.bookkeeper.stats.NullStatsProvider; +import org.apache.bookkeeper.stats.StatsLogger; +import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.commons.configuration.Configuration; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.stats.prometheus.metrics.PrometheusMetricsProvider; import org.apache.zookeeper.ZooKeeper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -45,6 +51,7 @@ public class ManagedLedgerClientFactory implements Closeable { private final ManagedLedgerFactory managedLedgerFactory; private final BookKeeper defaultBkClient; private final Map bkEnsemblePolicyToBkClientMap = Maps.newConcurrentMap(); + private StatsProvider statsProvider = new NullStatsProvider(); public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, BookKeeperClientFactory bookkeeperProvider) throws Exception { @@ -58,7 +65,17 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, managedLedgerFactoryConfig.setCopyEntriesInCache(conf.isManagedLedgerCacheCopyEntries()); managedLedgerFactoryConfig.setPrometheusStatsLatencyRolloverSeconds(conf.getManagedLedgerPrometheusStatsLatencyRolloverSeconds()); managedLedgerFactoryConfig.setTraceTaskExecution(conf.isManagedLedgerTraceTaskExecution()); - managedLedgerFactoryConfig.setClusterName(conf.getClusterName()); + + Configuration configuration = new ClientConfiguration(); + if (conf.isBookkeeperClientExposeStatsToPrometheus()) { + configuration.addProperty(PrometheusMetricsProvider.PROMETHEUS_STATS_LATENCY_ROLLOVER_SECONDS, + conf.getManagedLedgerPrometheusStatsLatencyRolloverSeconds()); + configuration.addProperty(PrometheusMetricsProvider.CLUSTER_NAME, conf.getClusterName()); + statsProvider = new PrometheusMetricsProvider(); + } + + statsProvider.start(configuration); + StatsLogger statsLogger = statsProvider.getStatsLogger("pulsar_managedLedger_client"); this.defaultBkClient = bookkeeperProvider.create(conf, zkClient, Optional.empty(), null); @@ -83,7 +100,7 @@ public ManagedLedgerClientFactory(ServiceConfiguration conf, ZooKeeper zkClient, return bkClient != null ? bkClient : defaultBkClient; }; - this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkFactory, zkClient, managedLedgerFactoryConfig); + this.managedLedgerFactory = new ManagedLedgerFactoryImpl(bkFactory, zkClient, managedLedgerFactoryConfig, statsLogger); } public ManagedLedgerFactory getManagedLedgerFactory() { @@ -94,16 +111,22 @@ public BookKeeper getBookKeeperClient() { return defaultBkClient; } + public StatsProvider getStatsProvider() { + return statsProvider; + } + @VisibleForTesting public Map getBkEnsemblePolicyToBookKeeperMap() { return bkEnsemblePolicyToBkClientMap; } + @Override public void close() throws IOException { try { managedLedgerFactory.shutdown(); log.info("Closed managed ledger factory"); + statsProvider.stop(); try { defaultBkClient.close(); } catch (RejectedExecutionException ree) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java index 15e28125a8cea..745eef8dd56cb 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/PrometheusMetricsGenerator.java @@ -28,6 +28,7 @@ import java.io.Writer; import java.util.Enumeration; +import org.apache.bookkeeper.stats.NullStatsProvider; import org.apache.bookkeeper.stats.StatsProvider; import org.apache.pulsar.broker.PulsarService; import static org.apache.pulsar.common.stats.JvmMetrics.getJvmDirectMemoryUsed; @@ -161,7 +162,11 @@ private static void parseMetricsToPrometheusMetrics(Collection metrics, } private static void generateManagedLedgerBookieClientMetrics(PulsarService pulsar, SimpleTextOutputStream stream) { - StatsProvider statsProvider = pulsar.getManagedLedgerClientFactory().getManagedLedgerFactory().getStatsProvider(); + StatsProvider statsProvider = pulsar.getManagedLedgerClientFactory().getStatsProvider(); + if (statsProvider instanceof NullStatsProvider) { + return; + } + try { Writer writer = new StringWriter(); statsProvider.writeAllMetrics(writer); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/DataSketchesOpStatsLogger.java similarity index 99% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/DataSketchesOpStatsLogger.java index dc5cdd76993df..3ef453ddc654e 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/DataSketchesOpStatsLogger.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/DataSketchesOpStatsLogger.java @@ -16,14 +16,15 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.bookkeeper.mledger.stats.prometheus; +package org.apache.pulsar.broker.stats.prometheus.metrics; import com.yahoo.sketches.quantiles.DoublesSketch; import com.yahoo.sketches.quantiles.DoublesSketchBuilder; import com.yahoo.sketches.quantiles.DoublesUnion; import com.yahoo.sketches.quantiles.DoublesUnionBuilder; - import io.netty.util.concurrent.FastThreadLocal; +import org.apache.bookkeeper.stats.OpStatsData; +import org.apache.bookkeeper.stats.OpStatsLogger; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -31,9 +32,6 @@ import java.util.concurrent.atomic.LongAdder; import java.util.concurrent.locks.StampedLock; -import org.apache.bookkeeper.stats.OpStatsData; -import org.apache.bookkeeper.stats.OpStatsLogger; - /** * OpStatsLogger implementation that uses DataSketches library to calculate the approximated latency quantiles. */ diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusMetricsProvider.java similarity index 94% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusMetricsProvider.java index dc8b980b116c0..a8a05e172502f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusMetricsProvider.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusMetricsProvider.java @@ -16,21 +16,11 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.bookkeeper.mledger.stats.prometheus; +package org.apache.pulsar.broker.stats.prometheus.metrics; import com.google.common.annotations.VisibleForTesting; - import io.netty.util.concurrent.DefaultThreadFactory; import io.prometheus.client.Collector; - -import java.io.IOException; -import java.io.Writer; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.ConcurrentSkipListMap; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; - import org.apache.bookkeeper.stats.CachingStatsProvider; import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.stats.StatsProvider; @@ -38,6 +28,10 @@ import org.apache.commons.configuration.Configuration; import org.apache.commons.lang.StringUtils; +import java.io.IOException; +import java.io.Writer; +import java.util.concurrent.*; + /** * A Prometheus based {@link StatsProvider} implementation. */ @@ -132,5 +126,4 @@ void rotateLatencyCollection() { metric.rotateLatencyCollection(); }); } - } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusStatsLogger.java similarity index 97% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusStatsLogger.java index ca8180d0b67b0..ad3c62f78d731 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusStatsLogger.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusStatsLogger.java @@ -16,17 +16,14 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.bookkeeper.mledger.stats.prometheus; +package org.apache.pulsar.broker.stats.prometheus.metrics; import com.google.common.base.Joiner; - import io.prometheus.client.Collector; - import org.apache.bookkeeper.stats.Counter; import org.apache.bookkeeper.stats.Gauge; import org.apache.bookkeeper.stats.OpStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; - import org.apache.bookkeeper.stats.prometheus.LongAdderCounter; /** diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusTextFormatUtil.java similarity index 87% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusTextFormatUtil.java index 30f45a601a134..abe0b560aa746 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/PrometheusTextFormatUtil.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/PrometheusTextFormatUtil.java @@ -16,19 +16,18 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.bookkeeper.mledger.stats.prometheus; +package org.apache.pulsar.broker.stats.prometheus.metrics; import io.prometheus.client.Collector; import io.prometheus.client.Collector.MetricFamilySamples; import io.prometheus.client.Collector.MetricFamilySamples.Sample; import io.prometheus.client.CollectorRegistry; +import org.apache.bookkeeper.stats.Counter; import java.io.IOException; import java.io.Writer; import java.util.Enumeration; -import org.apache.bookkeeper.stats.Counter; - /** * Logic to write metrics in Prometheus text format. */ @@ -36,7 +35,7 @@ public class PrometheusTextFormatUtil { static void writeGauge(Writer w, String name, String cluster, SimpleGauge gauge) { // Example: // # TYPE bookie_client_bookkeeper_ml_scheduler_completed_tasks_0 gauge - // pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0{cluster="test"} 1044057 + // pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0{cluster="pulsar"} 1044057 try { w.append("# TYPE ").append(name).append(" gauge\n"); w.append(name).append("{cluster=\"").append(cluster).append("\"}") @@ -62,24 +61,24 @@ static void writeCounter(Writer w, String name, String cluster, Counter counter) static void writeOpStat(Writer w, String name, String cluster, DataSketchesOpStatsLogger opStat) { // Example: // # TYPE pulsar_bookie_client_bookkeeper_ml_workers_task_queued summary - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="0.5"} NaN - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="0.75"} NaN - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="0.95"} NaN - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="0.99"} NaN - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="0.999"} NaN - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="0.9999"} NaN - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="false", quantile="1.0"} -Infinity - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_count{cluster="test", success="false"} 0 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{cluster="test", success="false"} 0.0 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="0.5"} 0.031 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="0.75"} 0.043 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="0.95"} 0.061 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="0.99"} 0.064 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="0.999"} 0.073 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="0.9999"} 0.073 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="test", success="true", quantile="1.0"} 0.552 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_count{cluster="test", success="true"} 40911432 - // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{cluster="test", success="true"} 527.0 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="0.5"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="0.75"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="0.95"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="0.99"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="0.999"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="0.9999"} NaN + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="false", quantile="1.0"} -Infinity + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_count{cluster="pulsar", success="false"} 0 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{cluster="pulsar", success="false"} 0.0 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="0.5"} 0.031 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="0.75"} 0.043 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="0.95"} 0.061 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="0.99"} 0.064 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="0.999"} 0.073 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="0.9999"} 0.073 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued{cluster="pulsar", success="true", quantile="1.0"} 0.552 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_count{cluster="pulsar", success="true"} 40911432 + // pulsar_bookie_client_bookkeeper_ml_workers_task_queued_sum{cluster="pulsar", success="true"} 527.0 try { w.append("# TYPE ").append(name).append(" summary\n"); writeQuantile(w, opStat, name, cluster,false, 0.5); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/SimpleGauge.java similarity index 95% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/SimpleGauge.java index f6065611e3cf2..a93a26c1bd1e6 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/stats/prometheus/SimpleGauge.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/stats/prometheus/metrics/SimpleGauge.java @@ -16,7 +16,7 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.bookkeeper.mledger.stats.prometheus; +package org.apache.pulsar.broker.stats.prometheus.metrics; import org.apache.bookkeeper.stats.Gauge; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/MockedBookKeeperClientFactory.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/MockedBookKeeperClientFactory.java index 4b7f49ba6f759..7c66ccc36c1d4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/MockedBookKeeperClientFactory.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/MockedBookKeeperClientFactory.java @@ -31,6 +31,7 @@ import org.apache.bookkeeper.client.EnsemblePlacementPolicy; import org.apache.bookkeeper.client.PulsarMockBookKeeper; +import org.apache.bookkeeper.stats.StatsLogger; import org.apache.zookeeper.ZooKeeper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -60,6 +61,13 @@ public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, return mockedBk; } + @Override + public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, + Optional> ensemblePlacementPolicyClass, + Map properties, StatsLogger statsLogger) throws IOException { + return mockedBk; + } + @Override public void close() { try { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index 53e629be4e386..034f1d7591b31 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -43,6 +43,7 @@ import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.EnsemblePlacementPolicy; import org.apache.bookkeeper.client.PulsarMockBookKeeper; +import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.util.ZkUtils; import org.apache.pulsar.broker.BookKeeperClientFactory; import org.apache.pulsar.broker.NoOpShutdownService; @@ -306,6 +307,14 @@ public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, return mockBookKeeper; } + @Override + public BookKeeper create(ServiceConfiguration conf, ZooKeeper zkClient, + Optional> ensemblePlacementPolicyClass, + Map properties, StatsLogger statsLogger) { + // Always return the same instance (so that we don't loose the mock BK content on broker restart + return mockBookKeeper; + } + @Override public void close() { // no-op diff --git a/site2/docs/reference-metrics.md b/site2/docs/reference-metrics.md index 03a4bb757eadb..73c012f71e762 100644 --- a/site2/docs/reference-metrics.md +++ b/site2/docs/reference-metrics.md @@ -328,16 +328,16 @@ All the managed ledger bookie client metrics labelled with the following labels: | Name | Type | Description | | --- | --- | --- | -| pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_* | Gauge | The number of tasks the scheduler executor execute completed.
The number of metrics determined by the scheduler executor thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| -| pulsar_bookie_client_bookkeeper_ml_scheduler_queue_* | Gauge | The number of tasks queued in the scheduler executor's queue.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| -| pulsar_bookie_client_bookkeeper_ml_scheduler_total_tasks_* | Gauge | The total number of tasks the scheduler executor received.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| -| pulsar_bookie_client_bookkeeper_ml_workers_completed_tasks_* | Gauge | The number of tasks the worker executor execute completed.
The number of metrics determined by the number of worker task thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`
| -| pulsar_bookie_client_bookkeeper_ml_workers_queue_* | Gauge | The number of tasks queued in the worker executor's queue.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`.
| -| pulsar_bookie_client_bookkeeper_ml_workers_total_tasks_* | Gauge | The total number of tasks the worker executor received.
The number of metrics determined by worker executor's thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`.
| -| pulsar_bookie_client_bookkeeper_ml_scheduler_task_execution | Summary | The scheduler task execution latency calculated in milliseconds. | -| pulsar_bookie_client_bookkeeper_ml_scheduler_task_queued | Summary | The scheduler task queued latency calculated in milliseconds. | -| pulsar_bookie_client_bookkeeper_ml_workers_task_execution | Summary | The worker task execution latency calculated in milliseconds. | -| pulsar_bookie_client_bookkeeper_ml_workers_task_queued | Summary | The worker task queued latency calculated in milliseconds. | +| pulsar_managedLedger_client_bookkeeper_ml_scheduler_completed_tasks_* | Gauge | The number of tasks the scheduler executor execute completed.
The number of metrics determined by the scheduler executor thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| +| pulsar_managedLedger_client_bookkeeper_ml_scheduler_queue_* | Gauge | The number of tasks queued in the scheduler executor's queue.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| +| pulsar_managedLedger_client_bookkeeper_ml_scheduler_total_tasks_* | Gauge | The total number of tasks the scheduler executor received.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumSchedulerThreads` in `broker.conf`.
| +| pulsar_managedLedger_client_bookkeeper_ml_workers_completed_tasks_* | Gauge | The number of tasks the worker executor execute completed.
The number of metrics determined by the number of worker task thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`
| +| pulsar_managedLedger_client_bookkeeper_ml_workers_queue_* | Gauge | The number of tasks queued in the worker executor's queue.
The number of metrics determined by scheduler executor's thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`.
| +| pulsar_managedLedger_client_bookkeeper_ml_workers_total_tasks_* | Gauge | The total number of tasks the worker executor received.
The number of metrics determined by worker executor's thread number configured by `managedLedgerNumWorkerThreads` in `broker.conf`.
| +| pulsar_managedLedger_client_bookkeeper_ml_scheduler_task_execution | Summary | The scheduler task execution latency calculated in milliseconds. | +| pulsar_managedLedger_client_bookkeeper_ml_scheduler_task_queued | Summary | The scheduler task queued latency calculated in milliseconds. | +| pulsar_managedLedger_client_bookkeeper_ml_workers_task_execution | Summary | The worker task execution latency calculated in milliseconds. | +| pulsar_managedLedger_client_bookkeeper_ml_workers_task_queued | Summary | The worker task queued latency calculated in milliseconds. | ## Monitor From dfabeeccfcb62bdc0dab20ef312160c721295133 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Mon, 4 May 2020 20:25:37 +0800 Subject: [PATCH 10/12] remove unnecessary dependency --- .../bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index f590e1e27bc2a..10b3146a1b9d7 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -73,13 +73,10 @@ import org.apache.bookkeeper.mledger.proto.MLDataFormats.LongProperty; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedCursorInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.MessageRange; -import org.apache.bookkeeper.mledger.stats.prometheus.PrometheusMetricsProvider; import org.apache.bookkeeper.mledger.util.Futures; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; -import org.apache.bookkeeper.stats.StatsProvider; import org.apache.bookkeeper.zookeeper.ZooKeeperClient; -import org.apache.commons.configuration.Configuration; import org.apache.pulsar.common.util.DateFormatter; import org.apache.pulsar.metadata.api.MetadataStore; import org.apache.pulsar.metadata.api.Stat; From b79789f75bbce88585106468d088220897db9313 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Wed, 6 May 2020 11:29:58 +0800 Subject: [PATCH 11/12] fix test case --- .../pulsar/broker/stats/PrometheusMetricsTest.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java index a1472e7ba4833..487c971f370d2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/PrometheusMetricsTest.java @@ -329,23 +329,23 @@ public void testManagedLedgerBookieClientStats() throws Exception { System.out.println(e.getKey() + ": " + e.getValue()) ); - List cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_scheduler_completed_tasks_0"); + List cm = (List) metrics.get("pulsar_managedLedger_client_bookkeeper_ml_scheduler_completed_tasks_0"); assertEquals(cm.size(), 1); assertEquals(cm.get(0).tags.get("cluster"), "test"); - cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_scheduler_queue_0"); + cm = (List) metrics.get("pulsar_managedLedger_client_bookkeeper_ml_scheduler_queue_0"); assertEquals(cm.size(), 1); assertEquals(cm.get(0).tags.get("cluster"), "test"); - cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_scheduler_total_tasks_0"); + cm = (List) metrics.get("pulsar_managedLedger_client_bookkeeper_ml_scheduler_total_tasks_0"); assertEquals(cm.size(), 1); assertEquals(cm.get(0).tags.get("cluster"), "test"); - cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_workers_completed_tasks_0"); + cm = (List) metrics.get("pulsar_managedLedger_client_bookkeeper_ml_workers_completed_tasks_0"); assertEquals(cm.size(), 1); assertEquals(cm.get(0).tags.get("cluster"), "test"); - cm = (List) metrics.get("pulsar_bookie_client_bookkeeper_ml_workers_task_execution_count"); + cm = (List) metrics.get("pulsar_managedLedger_client_bookkeeper_ml_workers_task_execution_count"); assertEquals(cm.size(), 2); assertEquals(cm.get(0).tags.get("cluster"), "test"); From 1923536294df6d3e3466cb3f28e34aca375b6bcd Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Wed, 6 May 2020 14:08:04 +0800 Subject: [PATCH 12/12] fix test case --- .../apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index 034f1d7591b31..538a36a5cfb32 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -20,6 +20,7 @@ import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; import com.google.common.collect.Sets; import com.google.common.util.concurrent.MoreExecutors; @@ -107,6 +108,7 @@ protected void resetConfig() { this.conf.setBrokerServicePortTls(Optional.of(0)); this.conf.setWebServicePort(Optional.of(0)); this.conf.setWebServicePortTls(Optional.of(0)); + this.conf.setBookkeeperClientExposeStatsToPrometheus(true); } protected final void internalSetup() throws Exception {