From 18c7584432016e07e801ebbc75ca304bc8324e13 Mon Sep 17 00:00:00 2001 From: daojun Date: Fri, 12 Aug 2022 02:55:37 +0800 Subject: [PATCH 01/10] add BatchMetadataStoreStats --- .../broker/stats/PrometheusMetricsTest.java | 86 ++++++++++- pulsar-metadata/pom.xml | 5 + .../metadata/impl/AbstractMetadataStore.java | 5 +- .../AbstractBatchedMetadataStore.java | 26 +++- .../metadata/impl/batching/MetadataOp.java | 2 + .../metadata/impl/batching/OpDelete.java | 6 + .../pulsar/metadata/impl/batching/OpGet.java | 6 + .../metadata/impl/batching/OpGetChildren.java | 6 + .../pulsar/metadata/impl/batching/OpPut.java | 6 + .../impl/stats/BatchMetadataStoreStats.java | 135 ++++++++++++++++++ .../metadata/impl/stats/package-info.java | 19 +++ 11 files changed, 288 insertions(+), 14 deletions(-) create mode 100644 pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java create mode 100644 pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/package-info.java 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 c5ecb8d5bf6a6..3f48be8c4d7f2 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 @@ -65,12 +65,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsGenerator; -import org.apache.pulsar.client.api.Consumer; -import org.apache.pulsar.client.api.MessageRoutingMode; -import org.apache.pulsar.client.api.Producer; -import org.apache.pulsar.client.api.PulsarClient; -import org.apache.pulsar.client.api.PulsarClientException; -import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.client.api.*; import org.apache.pulsar.compaction.Compactor; import org.awaitility.Awaitility; import org.mockito.Mockito; @@ -1522,6 +1517,85 @@ public void testSplitTopicAndPartitionLabel() throws Exception { consumer2.close(); } + + @Test + public void testBatchMetadataStoreMetrics() throws Exception { + String ns = "prop/ns-abc1"; + admin.namespaces().createNamespace(ns); + + String topic = "persistent://prop/ns-abc1/metadata-store-" + UUID.randomUUID(); + String subName = "my-sub1"; + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic).create(); + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic).subscriptionName(subName).subscribe(); + + for (int i = 0; i < 100; i++) { + producer.newMessage().value(UUID.randomUUID().toString()).send(); + } + + for (;;) { + Message message = consumer.receive(10, TimeUnit.SECONDS); + if (message == null) { + break; + } + consumer.acknowledge(message); + } + + ByteArrayOutputStream output = new ByteArrayOutputStream(); + PrometheusMetricsGenerator.generate(pulsar, false, false, false, false, output); + String metricsStr = output.toString(); + Multimap metricsMap = parseMetrics(metricsStr); + + Collection readOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_read_ops_overflow" + "_total"); + Collection writeOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_write_ops_overflow" + "_total"); + Collection queueingWriteOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write_ops"); + Collection queueingReadOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write_ops"); + Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); + Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_op_waiting_ms" + "_sum"); + + Assert.assertTrue(readOpsOverflow.size() > 1); + Assert.assertTrue(writeOpsOverflow.size() > 1); + Assert.assertTrue(queueingWriteOps.size() > 1); + Assert.assertTrue(queueingReadOps.size() > 1); + Assert.assertTrue(executorQueueSize.size() > 1); + Assert.assertTrue(opsWaiting.size() > 1); + + for (Metric m : readOpsOverflow) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.value >= 0); + } + for (Metric m : writeOpsOverflow) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.value >= 0); + } + for (Metric m : queueingWriteOps) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.value >= 0); + } + for (Metric m : queueingReadOps) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.value >= 0); + } + for (Metric m : executorQueueSize) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.value >= 0); + } + for (Metric m : opsWaiting) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.value >= 0); + } + } + private void compareCompactionStateCount(List cm, double count) { assertEquals(cm.size(), 1); assertEquals(cm.get(0).tags.get("cluster"), "test"); diff --git a/pulsar-metadata/pom.xml b/pulsar-metadata/pom.xml index 935fc878a628c..b6d33593abf59 100644 --- a/pulsar-metadata/pom.xml +++ b/pulsar-metadata/pom.xml @@ -99,6 +99,11 @@ caffeine + + io.prometheus + simpleclient + + diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java index 8986811ad9f71..358e449a76e34 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/AbstractMetadataStore.java @@ -35,9 +35,9 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.Executor; -import java.util.concurrent.Executors; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import java.util.stream.Collectors; @@ -81,8 +81,7 @@ public abstract class AbstractMetadataStore implements MetadataStoreExtended, Co protected abstract CompletableFuture existsFromStore(String path); protected AbstractMetadataStore() { - this.executor = Executors - .newSingleThreadScheduledExecutor(new DefaultThreadFactory("metadata-store")); + this.executor = new ScheduledThreadPoolExecutor(1, new DefaultThreadFactory("metadata-store")); registerListener(this); this.childrenCache = Caffeine.newBuilder() diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index c9d245b8caf46..cec947bf075b1 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -34,6 +34,7 @@ import org.apache.pulsar.metadata.api.Stat; import org.apache.pulsar.metadata.api.extended.CreateOption; import org.apache.pulsar.metadata.impl.AbstractMetadataStore; +import org.apache.pulsar.metadata.impl.stats.BatchMetadataStoreStats; import org.jctools.queues.MessagePassingQueue; import org.jctools.queues.MpscUnboundedArrayQueue; @@ -50,7 +51,8 @@ public abstract class AbstractBatchedMetadataStore extends AbstractMetadataStore private final int maxDelayMillis; private final int maxOperations; private final int maxSize; - private MetadataEventSynchronizer synchronizer; + private final MetadataEventSynchronizer synchronizer; + private final BatchMetadataStoreStats batchMetadataStoreStats; protected AbstractBatchedMetadataStore(MetadataStoreConfig conf) { super(); @@ -74,6 +76,7 @@ protected AbstractBatchedMetadataStore(MetadataStoreConfig conf) { // update synchronizer and register sync listener synchronizer = conf.getSynchronizer(); registerSyncLister(Optional.ofNullable(synchronizer)); + this.batchMetadataStoreStats = new BatchMetadataStoreStats(executor, this.readOps, this.writeOps); } @Override @@ -93,7 +96,7 @@ private void flush() { while (!readOps.isEmpty()) { List ops = new ArrayList<>(); readOps.drain(ops::add, maxOperations); - batchOperation(ops); + internalBatchOperation(ops); } while (!writeOps.isEmpty()) { @@ -114,7 +117,7 @@ private void flush() { batchSize += op.size(); ops.add(writeOps.poll()); } - batchOperation(ops); + internalBatchOperation(ops); } flushInProgress.set(false); @@ -157,16 +160,29 @@ public Optional getMetadataEventSynchronizer() { private void enqueue(MessagePassingQueue queue, MetadataOp op) { if (enabled) { if (!queue.offer(op)) { + if (queue.equals(this.readOps)) { + this.batchMetadataStoreStats.recordReadOpOverflow(); + } else { + this.batchMetadataStoreStats.recordWriteOpOverflow(); + } // Execute individually if we're failing to enqueue - batchOperation(Collections.singletonList(op)); + internalBatchOperation(Collections.singletonList(op)); return; } if (queue.size() > maxOperations && flushInProgress.compareAndSet(false, true)) { executor.execute(this::flush); } } else { - batchOperation(Collections.singletonList(op)); + internalBatchOperation(Collections.singletonList(op)); + } + } + + private void internalBatchOperation(List ops) { + long now = System.currentTimeMillis(); + for (MetadataOp op : ops) { + this.batchMetadataStoreStats.recordOpWaiting(now - op.created()); } + this.batchOperation(ops); } protected abstract void batchOperation(List ops); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/MetadataOp.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/MetadataOp.java index 6c78c587aedd9..599e525ef3801 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/MetadataOp.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/MetadataOp.java @@ -34,6 +34,8 @@ enum Type { int size(); + long created(); + default OpGet asGet() { return (OpGet) this; } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpDelete.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpDelete.java index 6a039c302eb32..6de8f9320575e 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpDelete.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpDelete.java @@ -28,6 +28,7 @@ public class OpDelete implements MetadataOp { private final String path; private final Optional optExpectedVersion; + public final long created = System.currentTimeMillis(); private final CompletableFuture future = new CompletableFuture<>(); @@ -40,4 +41,9 @@ public Type getType() { public int size() { return path.length(); } + + @Override + public long created() { + return this.created; + } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGet.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGet.java index 5d9c6e3536a1b..6c4832e479c2b 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGet.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGet.java @@ -29,6 +29,7 @@ public class OpGet implements MetadataOp { private final String path; + public final long created = System.currentTimeMillis(); private final CompletableFuture> future = new CompletableFuture<>(); @Override @@ -40,4 +41,9 @@ public Type getType() { public int size() { return path.length(); } + + @Override + public long created() { + return this.created; + } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGetChildren.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGetChildren.java index 00220f4ba4895..4a435e7e350d1 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGetChildren.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpGetChildren.java @@ -28,6 +28,7 @@ public class OpGetChildren implements MetadataOp { private final String path; + public final long created = System.currentTimeMillis(); private final CompletableFuture> future = new CompletableFuture<>(); @Override @@ -39,4 +40,9 @@ public Type getType() { public int size() { return path.length(); } + + @Override + public long created() { + return this.created; + } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpPut.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpPut.java index cb3832ecfc4ac..36829971ae7ae 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpPut.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/OpPut.java @@ -34,6 +34,7 @@ public class OpPut implements MetadataOp { private final Optional optExpectedVersion; private final EnumSet options; + public final long created = System.currentTimeMillis(); private final CompletableFuture future = new CompletableFuture<>(); public boolean isEphemeral() { @@ -49,4 +50,9 @@ public Type getType() { public int size() { return path.length() + (data != null ? data.length : 0); } + + @Override + public long created() { + return this.created; + } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java new file mode 100644 index 0000000000000..31fc567df1c02 --- /dev/null +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -0,0 +1,135 @@ +/** + * 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.pulsar.metadata.impl.stats; + +import io.prometheus.client.Counter; +import io.prometheus.client.Gauge; +import io.prometheus.client.Histogram; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.pulsar.metadata.impl.batching.MetadataOp; +import org.jctools.queues.MessagePassingQueue; + +public final class BatchMetadataStoreStats implements AutoCloseable { + private static final double[] BUCKETS = new double[]{1, 5, 10, 20, 50, 100, 200, 500, 1000}; + private static final AtomicInteger COUNTER = new AtomicInteger(1); + private static final String LABEL_NAME = "metadata_store_name"; + + private static final Gauge QUEUEING_READ_OPS = Gauge + .build("pulsar_batch_metadata_store_queueing_read_ops", "-") + .labelNames(LABEL_NAME) + .register(); + private static final Gauge QUEUEING_WRITE_OPS = Gauge + .build("pulsar_batch_metadata_store_queueing_write_ops", "-") + .labelNames(LABEL_NAME) + .register(); + private static final Gauge EXECUTOR_QUEUE_SIZE = Gauge + .build("pulsar_batch_metadata_store_executor_queue_size", "-") + .labelNames(LABEL_NAME) + .register(); + private static final Counter READ_OPS_OVERFLOW = Counter + .build("pulsar_batch_metadata_store_read_ops_overflow" , "-") + .labelNames(LABEL_NAME) + .register(); + private static final Counter WRITE_OPS_OVERFLOW = Counter + .build("pulsar_batch_metadata_store_write_ops_overflow" , "-") + .labelNames(LABEL_NAME) + .register(); + private static final Histogram BATCH_OPS_WAITING = Histogram + .build("pulsar_batch_metadata_store_op_waiting", "-") + .labelNames(LABEL_NAME) + .unit("ms") + .buckets(BUCKETS) + .register(); + + private final AtomicBoolean closed = new AtomicBoolean(false); + private final ThreadPoolExecutor executor; + private final MessagePassingQueue readOps; + private final MessagePassingQueue writeOps; + private final String metadataStoreName; + + private final Counter.Child readOpsOverflowChild; + private final Counter.Child writeOpsOverflowChild; + private final Histogram.Child batchOpsWaitingChild; + + public BatchMetadataStoreStats(ExecutorService executor, MessagePassingQueue readOps, + MessagePassingQueue writeOps) { + if (executor instanceof ThreadPoolExecutor tx) { + this.executor = tx; + } else { + this.executor = null; + } + this.readOps = readOps; + this.writeOps = writeOps; + this.metadataStoreName = "metadata_store_" + COUNTER.getAndIncrement(); + + + QUEUEING_READ_OPS.setChild(new Gauge.Child() { + @Override + public double get() { + return BatchMetadataStoreStats.this.readOps.size(); + } + }, metadataStoreName); + + QUEUEING_WRITE_OPS.setChild(new Gauge.Child() { + @Override + public double get() { + return BatchMetadataStoreStats.this.writeOps.size(); + } + }, metadataStoreName); + + EXECUTOR_QUEUE_SIZE.setChild(new Gauge.Child() { + @Override + public double get() { + return BatchMetadataStoreStats.this.executor == null ? 0 : + BatchMetadataStoreStats.this.executor.getQueue().size(); + } + }, metadataStoreName); + + this.readOpsOverflowChild = READ_OPS_OVERFLOW.labels(metadataStoreName); + this.writeOpsOverflowChild = WRITE_OPS_OVERFLOW.labels(metadataStoreName); + this.batchOpsWaitingChild = BATCH_OPS_WAITING.labels(metadataStoreName); + } + + public void recordOpWaiting(long millis) { + this.batchOpsWaitingChild.observe(millis); + } + + public void recordReadOpOverflow() { + this.readOpsOverflowChild.inc(); + } + + public void recordWriteOpOverflow() { + this.writeOpsOverflowChild.inc(); + } + + @Override + public void close() throws Exception { + if (closed.compareAndSet(false, true)) { + QUEUEING_READ_OPS.remove(this.metadataStoreName); + QUEUEING_WRITE_OPS.remove(this.metadataStoreName); + EXECUTOR_QUEUE_SIZE.remove(this.metadataStoreName); + BATCH_OPS_WAITING.remove(this.metadataStoreName); + READ_OPS_OVERFLOW.remove(this.metadataStoreName); + WRITE_OPS_OVERFLOW.remove(this.metadataStoreName); + } + } +} diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/package-info.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/package-info.java new file mode 100644 index 0000000000000..15ca0d1c58263 --- /dev/null +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/package-info.java @@ -0,0 +1,19 @@ +/** + * 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.pulsar.metadata.impl.stats; From 3e06fefc63b37d9cda3de9feaecc2e1b839f5335 Mon Sep 17 00:00:00 2001 From: daojun Date: Fri, 12 Aug 2022 03:43:18 +0800 Subject: [PATCH 02/10] upate stats --- .../metadata/impl/batching/AbstractBatchedMetadataStore.java | 1 + 1 file changed, 1 insertion(+) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index cec947bf075b1..10293bce842be 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -90,6 +90,7 @@ public void close() throws Exception { scheduledTask.cancel(true); } super.close(); + this.batchMetadataStoreStats.close(); } private void flush() { From 51a23613689e54d3224b57677c35bebca4a0d1c7 Mon Sep 17 00:00:00 2001 From: daojun Date: Fri, 12 Aug 2022 11:47:28 +0800 Subject: [PATCH 03/10] review fix --- .../broker/stats/PrometheusMetricsTest.java | 31 ++++++++++++------- .../impl/stats/BatchMetadataStoreStats.java | 14 ++++----- 2 files changed, 26 insertions(+), 19 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 3f48be8c4d7f2..a408f9e846567 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 @@ -65,7 +65,14 @@ import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsGenerator; -import org.apache.pulsar.client.api.*; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.MessageRoutingMode; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.compaction.Compactor; import org.awaitility.Awaitility; import org.mockito.Mockito; @@ -1550,12 +1557,12 @@ public void testBatchMetadataStoreMetrics() throws Exception { String metricsStr = output.toString(); Multimap metricsMap = parseMetrics(metricsStr); - Collection readOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_read_ops_overflow" + "_total"); - Collection writeOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_write_ops_overflow" + "_total"); - Collection queueingWriteOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write_ops"); - Collection queueingReadOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write_ops"); + Collection readOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_read_overflow" + "_total"); + Collection writeOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_write_overflow" + "_total"); + Collection queueingWriteOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write"); + Collection queueingReadOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write"); Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); - Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_op_waiting_ms" + "_sum"); + Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_waiting_ms" + "_sum"); Assert.assertTrue(readOpsOverflow.size() > 1); Assert.assertTrue(writeOpsOverflow.size() > 1); @@ -1566,32 +1573,32 @@ public void testBatchMetadataStoreMetrics() throws Exception { for (Metric m : readOpsOverflow) { Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); Assert.assertTrue(m.value >= 0); } for (Metric m : writeOpsOverflow) { Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); Assert.assertTrue(m.value >= 0); } for (Metric m : queueingWriteOps) { Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); Assert.assertTrue(m.value >= 0); } for (Metric m : queueingReadOps) { Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); Assert.assertTrue(m.value >= 0); } for (Metric m : executorQueueSize) { Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); Assert.assertTrue(m.value >= 0); } for (Metric m : opsWaiting) { Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("metadata_store_name").startsWith("metadata_store_")); + Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); Assert.assertTrue(m.value >= 0); } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index 31fc567df1c02..650a206b172ce 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -31,14 +31,14 @@ public final class BatchMetadataStoreStats implements AutoCloseable { private static final double[] BUCKETS = new double[]{1, 5, 10, 20, 50, 100, 200, 500, 1000}; private static final AtomicInteger COUNTER = new AtomicInteger(1); - private static final String LABEL_NAME = "metadata_store_name"; + private static final String LABEL_NAME = "name"; private static final Gauge QUEUEING_READ_OPS = Gauge - .build("pulsar_batch_metadata_store_queueing_read_ops", "-") + .build("pulsar_batch_metadata_store_queueing_read", "-") .labelNames(LABEL_NAME) .register(); private static final Gauge QUEUEING_WRITE_OPS = Gauge - .build("pulsar_batch_metadata_store_queueing_write_ops", "-") + .build("pulsar_batch_metadata_store_queueing_write", "-") .labelNames(LABEL_NAME) .register(); private static final Gauge EXECUTOR_QUEUE_SIZE = Gauge @@ -46,17 +46,17 @@ public final class BatchMetadataStoreStats implements AutoCloseable { .labelNames(LABEL_NAME) .register(); private static final Counter READ_OPS_OVERFLOW = Counter - .build("pulsar_batch_metadata_store_read_ops_overflow" , "-") + .build("pulsar_batch_metadata_store_read_overflow" , "-") .labelNames(LABEL_NAME) .register(); private static final Counter WRITE_OPS_OVERFLOW = Counter - .build("pulsar_batch_metadata_store_write_ops_overflow" , "-") + .build("pulsar_batch_metadata_store_write_overflow" , "-") .labelNames(LABEL_NAME) .register(); private static final Histogram BATCH_OPS_WAITING = Histogram - .build("pulsar_batch_metadata_store_op_waiting", "-") - .labelNames(LABEL_NAME) + .build("pulsar_batch_metadata_store_waiting", "-") .unit("ms") + .labelNames(LABEL_NAME) .buckets(BUCKETS) .register(); From 5d5759f4ca11727cdbb828cb1511bb68c2f6209d Mon Sep 17 00:00:00 2001 From: daojun Date: Thu, 22 Sep 2022 21:56:47 +0800 Subject: [PATCH 04/10] fix --- .../broker/stats/MetadataStoreStatsTest.java | 104 ++++++++++++++++++ .../broker/stats/PrometheusMetricsTest.java | 81 -------------- .../AbstractBatchedMetadataStore.java | 3 +- .../impl/stats/BatchMetadataStoreStats.java | 57 +++++----- 4 files changed, 131 insertions(+), 114 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java index eba134a2c8dfe..cd55cfb1ec5f7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java @@ -137,4 +137,108 @@ public void testMetadataStoreStats() throws Exception { } } + @Test + public void testBatchMetadataStoreMetrics() throws Exception { + String ns = "prop/ns-abc1"; + admin.namespaces().createNamespace(ns); + + String topic = "persistent://prop/ns-abc1/metadata-store-" + UUID.randomUUID(); + String subName = "my-sub1"; + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic).create(); + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic).subscriptionName(subName).subscribe(); + + for (int i = 0; i < 100; i++) { + producer.newMessage().value(UUID.randomUUID().toString()).send(); + } + + for (;;) { + Message message = consumer.receive(10, TimeUnit.SECONDS); + if (message == null) { + break; + } + consumer.acknowledge(message); + } + + ByteArrayOutputStream output = new ByteArrayOutputStream(); + PrometheusMetricsGenerator.generate(pulsar, false, false, false, false, output); + String metricsStr = output.toString(); + Multimap metricsMap = PrometheusMetricsTest.parseMetrics(metricsStr); + + Collection opsOverflow = metricsMap.get("pulsar_batch_metadata_store_overflow_ops" + "_total"); + Collection queueingOps = metricsMap.get("pulsar_batch_metadata_store_queueing_ops"); + Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); + Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_waiting_ms" + "_sum"); + + Assert.assertTrue(opsOverflow.size() > 0 && opsOverflow.size() % 2 == 0); + Assert.assertTrue(queueingOps.size() > 0 && queueingOps.size() % 2 == 0); + Assert.assertTrue(executorQueueSize.size() > 1); + Assert.assertTrue(opsWaiting.size() > 1); + + int readOpsOverflow = 0; + int writeOpsOverflow = 0; + for (PrometheusMetricsTest.Metric m : opsOverflow) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + String metadataStoreName = m.tags.get("name"); + Assert.assertNotNull(metadataStoreName); + Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); + String opType = m.tags.get("type"); + Assert.assertNotNull(opType); + if (opType.equals("read")) { + readOpsOverflow++; + } else if (opType.equals("write")){ + writeOpsOverflow++; + } + Assert.assertTrue(m.value >= 0); + } + Assert.assertEquals(readOpsOverflow, writeOpsOverflow); + Assert.assertTrue(readOpsOverflow > 0); + + int queueingReadOps = 0; + int queueingWriteOps = 0; + for (PrometheusMetricsTest.Metric m : queueingOps) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + String metadataStoreName = m.tags.get("name"); + Assert.assertNotNull(metadataStoreName); + Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); + String opType = m.tags.get("type"); + Assert.assertNotNull(opType); + if (opType.equals("read")) { + queueingReadOps++; + } else if (opType.equals("write")){ + queueingWriteOps++; + } + Assert.assertTrue(m.value >= 0); + } + Assert.assertEquals(queueingReadOps, queueingWriteOps); + Assert.assertTrue(queueingReadOps > 0); + + for (PrometheusMetricsTest.Metric m : executorQueueSize) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + String metadataStoreName = m.tags.get("name"); + Assert.assertNotNull(metadataStoreName); + Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); + Assert.assertTrue(m.value >= 0); + } + for (PrometheusMetricsTest.Metric m : opsWaiting) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + String metadataStoreName = m.tags.get("name"); + Assert.assertNotNull(metadataStoreName); + Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); + Assert.assertTrue(m.value >= 0); + } + } + } 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 fb268f43d732b..cfab570301873 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 @@ -67,12 +67,10 @@ import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsGenerator; import org.apache.pulsar.client.api.Consumer; -import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; -import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.compaction.Compactor; import org.awaitility.Awaitility; @@ -1524,85 +1522,6 @@ public void testSplitTopicAndPartitionLabel() throws Exception { consumer2.close(); } - - @Test - public void testBatchMetadataStoreMetrics() throws Exception { - String ns = "prop/ns-abc1"; - admin.namespaces().createNamespace(ns); - - String topic = "persistent://prop/ns-abc1/metadata-store-" + UUID.randomUUID(); - String subName = "my-sub1"; - - @Cleanup - Producer producer = pulsarClient.newProducer(Schema.STRING) - .topic(topic).create(); - @Cleanup - Consumer consumer = pulsarClient.newConsumer(Schema.STRING) - .topic(topic).subscriptionName(subName).subscribe(); - - for (int i = 0; i < 100; i++) { - producer.newMessage().value(UUID.randomUUID().toString()).send(); - } - - for (;;) { - Message message = consumer.receive(10, TimeUnit.SECONDS); - if (message == null) { - break; - } - consumer.acknowledge(message); - } - - ByteArrayOutputStream output = new ByteArrayOutputStream(); - PrometheusMetricsGenerator.generate(pulsar, false, false, false, false, output); - String metricsStr = output.toString(); - Multimap metricsMap = parseMetrics(metricsStr); - - Collection readOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_read_overflow" + "_total"); - Collection writeOpsOverflow = metricsMap.get("pulsar_batch_metadata_store_write_overflow" + "_total"); - Collection queueingWriteOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write"); - Collection queueingReadOps = metricsMap.get("pulsar_batch_metadata_store_queueing_write"); - Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); - Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_waiting_ms" + "_sum"); - - Assert.assertTrue(readOpsOverflow.size() > 1); - Assert.assertTrue(writeOpsOverflow.size() > 1); - Assert.assertTrue(queueingWriteOps.size() > 1); - Assert.assertTrue(queueingReadOps.size() > 1); - Assert.assertTrue(executorQueueSize.size() > 1); - Assert.assertTrue(opsWaiting.size() > 1); - - for (Metric m : readOpsOverflow) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); - Assert.assertTrue(m.value >= 0); - } - for (Metric m : writeOpsOverflow) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); - Assert.assertTrue(m.value >= 0); - } - for (Metric m : queueingWriteOps) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); - Assert.assertTrue(m.value >= 0); - } - for (Metric m : queueingReadOps) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); - Assert.assertTrue(m.value >= 0); - } - for (Metric m : executorQueueSize) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); - Assert.assertTrue(m.value >= 0); - } - for (Metric m : opsWaiting) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - Assert.assertTrue(m.tags.get("name").startsWith("metadata_store_")); - Assert.assertTrue(m.value >= 0); - } - } - private void compareCompactionStateCount(List cm, double count) { assertEquals(cm.size(), 1); assertEquals(cm.get(0).tags.get("cluster"), "test"); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index 7603cc04cfea4..7d96cbf7e6d74 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -76,7 +76,8 @@ protected AbstractBatchedMetadataStore(MetadataStoreConfig conf) { // update synchronizer and register sync listener synchronizer = conf.getSynchronizer(); registerSyncLister(Optional.ofNullable(synchronizer)); - this.batchMetadataStoreStats = new BatchMetadataStoreStats(executor, this.readOps, this.writeOps); + this.batchMetadataStoreStats = + new BatchMetadataStoreStats(metadataStoreName, executor, this.readOps, this.writeOps); } @Override diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index 650a206b172ce..d8bc4edb12a4b 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -24,39 +24,32 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicInteger; import org.apache.pulsar.metadata.impl.batching.MetadataOp; import org.jctools.queues.MessagePassingQueue; public final class BatchMetadataStoreStats implements AutoCloseable { private static final double[] BUCKETS = new double[]{1, 5, 10, 20, 50, 100, 200, 500, 1000}; - private static final AtomicInteger COUNTER = new AtomicInteger(1); - private static final String LABEL_NAME = "name"; + private static final String NAME = "name"; + private static final String TYPE = "type"; + private static final String OP_READ = "read"; + private static final String OP_WRITE = "write"; - private static final Gauge QUEUEING_READ_OPS = Gauge - .build("pulsar_batch_metadata_store_queueing_read", "-") - .labelNames(LABEL_NAME) - .register(); - private static final Gauge QUEUEING_WRITE_OPS = Gauge - .build("pulsar_batch_metadata_store_queueing_write", "-") - .labelNames(LABEL_NAME) + private static final Gauge QUEUEING_OPS = Gauge + .build("pulsar_batch_metadata_store_queueing_ops", "-") + .labelNames(NAME, TYPE) .register(); private static final Gauge EXECUTOR_QUEUE_SIZE = Gauge .build("pulsar_batch_metadata_store_executor_queue_size", "-") - .labelNames(LABEL_NAME) - .register(); - private static final Counter READ_OPS_OVERFLOW = Counter - .build("pulsar_batch_metadata_store_read_overflow" , "-") - .labelNames(LABEL_NAME) + .labelNames(NAME) .register(); - private static final Counter WRITE_OPS_OVERFLOW = Counter - .build("pulsar_batch_metadata_store_write_overflow" , "-") - .labelNames(LABEL_NAME) + private static final Counter OVERFLOW_OPS = Counter + .build("pulsar_batch_metadata_store_overflow_ops" , "-") + .labelNames(NAME, TYPE) .register(); private static final Histogram BATCH_OPS_WAITING = Histogram .build("pulsar_batch_metadata_store_waiting", "-") .unit("ms") - .labelNames(LABEL_NAME) + .labelNames(NAME) .buckets(BUCKETS) .register(); @@ -70,8 +63,8 @@ public final class BatchMetadataStoreStats implements AutoCloseable { private final Counter.Child writeOpsOverflowChild; private final Histogram.Child batchOpsWaitingChild; - public BatchMetadataStoreStats(ExecutorService executor, MessagePassingQueue readOps, - MessagePassingQueue writeOps) { + public BatchMetadataStoreStats(String metadataStoreName, ExecutorService executor, + MessagePassingQueue readOps, MessagePassingQueue writeOps) { if (executor instanceof ThreadPoolExecutor tx) { this.executor = tx; } else { @@ -79,22 +72,22 @@ public BatchMetadataStoreStats(ExecutorService executor, MessagePassingQueue Date: Mon, 24 Oct 2022 14:09:21 +0800 Subject: [PATCH 05/10] review fix --- .../broker/stats/MetadataStoreStatsTest.java | 2 +- .../AbstractBatchedMetadataStore.java | 19 ++++++++++++------- .../impl/stats/BatchMetadataStoreStats.java | 2 +- 3 files changed, 14 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java index cd55cfb1ec5f7..f44f0856a696f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java @@ -172,7 +172,7 @@ public void testBatchMetadataStoreMetrics() throws Exception { Collection opsOverflow = metricsMap.get("pulsar_batch_metadata_store_overflow_ops" + "_total"); Collection queueingOps = metricsMap.get("pulsar_batch_metadata_store_queueing_ops"); Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); - Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_waiting_ms" + "_sum"); + Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_queue_wait_time_ms" + "_sum"); Assert.assertTrue(opsOverflow.size() > 0 && opsOverflow.size() % 2 == 0); Assert.assertTrue(queueingOps.size() > 0 && queueingOps.size() % 2 == 0); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index 7d96cbf7e6d74..307e5282a2674 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -128,21 +128,21 @@ private void flush() { @Override public final CompletableFuture> storeGet(String path) { OpGet op = new OpGet(path); - enqueue(readOps, op); + enqueue(readOps, op, OpType.READ); return op.getFuture(); } @Override protected final CompletableFuture> getChildrenFromStore(String path) { OpGetChildren op = new OpGetChildren(path); - enqueue(readOps, op); + enqueue(readOps, op, OpType.READ); return op.getFuture(); } @Override protected final CompletableFuture storeDelete(String path, Optional expectedVersion) { OpDelete op = new OpDelete(path, expectedVersion); - enqueue(writeOps, op); + enqueue(writeOps, op, OpType.WRITE); return op.getFuture(); } @@ -150,7 +150,7 @@ protected final CompletableFuture storeDelete(String path, Optional protected CompletableFuture storePut(String path, byte[] data, Optional optExpectedVersion, EnumSet options) { OpPut op = new OpPut(path, data, optExpectedVersion, options); - enqueue(writeOps, op); + enqueue(writeOps, op, OpType.WRITE); return op.getFuture(); } @@ -159,12 +159,12 @@ public Optional getMetadataEventSynchronizer() { return Optional.ofNullable(synchronizer); } - private void enqueue(MessagePassingQueue queue, MetadataOp op) { + private void enqueue(MessagePassingQueue queue, MetadataOp op, OpType opType) { if (enabled) { if (!queue.offer(op)) { - if (queue.equals(this.readOps)) { + if (opType.equals(OpType.READ)) { this.batchMetadataStoreStats.recordReadOpOverflow(); - } else { + } else if (opType.equals(OpType.WRITE)) { this.batchMetadataStoreStats.recordWriteOpOverflow(); } // Execute individually if we're failing to enqueue @@ -187,5 +187,10 @@ private void internalBatchOperation(List ops) { this.batchOperation(ops); } + private enum OpType { + READ, + WRITE + } + protected abstract void batchOperation(List ops); } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index d8bc4edb12a4b..dc74a640b22a1 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -47,7 +47,7 @@ public final class BatchMetadataStoreStats implements AutoCloseable { .labelNames(NAME, TYPE) .register(); private static final Histogram BATCH_OPS_WAITING = Histogram - .build("pulsar_batch_metadata_store_waiting", "-") + .build("pulsar_batch_metadata_store_queue_wait_time", "-") .unit("ms") .labelNames(NAME) .buckets(BUCKETS) From 869a0d963e84060cefa2e7c7e1e8a23dfa73ad03 Mon Sep 17 00:00:00 2001 From: daojun Date: Sat, 29 Oct 2022 21:03:23 +0800 Subject: [PATCH 06/10] fix --- .../broker/stats/MetadataStoreStatsTest.java | 46 ---------------- .../AbstractBatchedMetadataStore.java | 22 +++----- .../impl/stats/BatchMetadataStoreStats.java | 52 +------------------ 3 files changed, 7 insertions(+), 113 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java index f44f0856a696f..4b47be672d6a9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java @@ -169,58 +169,12 @@ public void testBatchMetadataStoreMetrics() throws Exception { String metricsStr = output.toString(); Multimap metricsMap = PrometheusMetricsTest.parseMetrics(metricsStr); - Collection opsOverflow = metricsMap.get("pulsar_batch_metadata_store_overflow_ops" + "_total"); - Collection queueingOps = metricsMap.get("pulsar_batch_metadata_store_queueing_ops"); Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_queue_wait_time_ms" + "_sum"); - Assert.assertTrue(opsOverflow.size() > 0 && opsOverflow.size() % 2 == 0); - Assert.assertTrue(queueingOps.size() > 0 && queueingOps.size() % 2 == 0); Assert.assertTrue(executorQueueSize.size() > 1); Assert.assertTrue(opsWaiting.size() > 1); - int readOpsOverflow = 0; - int writeOpsOverflow = 0; - for (PrometheusMetricsTest.Metric m : opsOverflow) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - String metadataStoreName = m.tags.get("name"); - Assert.assertNotNull(metadataStoreName); - Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) - || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) - || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); - String opType = m.tags.get("type"); - Assert.assertNotNull(opType); - if (opType.equals("read")) { - readOpsOverflow++; - } else if (opType.equals("write")){ - writeOpsOverflow++; - } - Assert.assertTrue(m.value >= 0); - } - Assert.assertEquals(readOpsOverflow, writeOpsOverflow); - Assert.assertTrue(readOpsOverflow > 0); - - int queueingReadOps = 0; - int queueingWriteOps = 0; - for (PrometheusMetricsTest.Metric m : queueingOps) { - Assert.assertEquals(m.tags.get("cluster"), "test"); - String metadataStoreName = m.tags.get("name"); - Assert.assertNotNull(metadataStoreName); - Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) - || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) - || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); - String opType = m.tags.get("type"); - Assert.assertNotNull(opType); - if (opType.equals("read")) { - queueingReadOps++; - } else if (opType.equals("write")){ - queueingWriteOps++; - } - Assert.assertTrue(m.value >= 0); - } - Assert.assertEquals(queueingReadOps, queueingWriteOps); - Assert.assertTrue(queueingReadOps > 0); - for (PrometheusMetricsTest.Metric m : executorQueueSize) { Assert.assertEquals(m.tags.get("cluster"), "test"); String metadataStoreName = m.tags.get("name"); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index 307e5282a2674..118b553bac985 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -77,7 +77,7 @@ protected AbstractBatchedMetadataStore(MetadataStoreConfig conf) { synchronizer = conf.getSynchronizer(); registerSyncLister(Optional.ofNullable(synchronizer)); this.batchMetadataStoreStats = - new BatchMetadataStoreStats(metadataStoreName, executor, this.readOps, this.writeOps); + new BatchMetadataStoreStats(metadataStoreName, executor); } @Override @@ -128,21 +128,21 @@ private void flush() { @Override public final CompletableFuture> storeGet(String path) { OpGet op = new OpGet(path); - enqueue(readOps, op, OpType.READ); + enqueue(readOps, op); return op.getFuture(); } @Override protected final CompletableFuture> getChildrenFromStore(String path) { OpGetChildren op = new OpGetChildren(path); - enqueue(readOps, op, OpType.READ); + enqueue(readOps, op); return op.getFuture(); } @Override protected final CompletableFuture storeDelete(String path, Optional expectedVersion) { OpDelete op = new OpDelete(path, expectedVersion); - enqueue(writeOps, op, OpType.WRITE); + enqueue(writeOps, op); return op.getFuture(); } @@ -150,7 +150,7 @@ protected final CompletableFuture storeDelete(String path, Optional protected CompletableFuture storePut(String path, byte[] data, Optional optExpectedVersion, EnumSet options) { OpPut op = new OpPut(path, data, optExpectedVersion, options); - enqueue(writeOps, op, OpType.WRITE); + enqueue(writeOps, op); return op.getFuture(); } @@ -159,14 +159,9 @@ public Optional getMetadataEventSynchronizer() { return Optional.ofNullable(synchronizer); } - private void enqueue(MessagePassingQueue queue, MetadataOp op, OpType opType) { + private void enqueue(MessagePassingQueue queue, MetadataOp op) { if (enabled) { if (!queue.offer(op)) { - if (opType.equals(OpType.READ)) { - this.batchMetadataStoreStats.recordReadOpOverflow(); - } else if (opType.equals(OpType.WRITE)) { - this.batchMetadataStoreStats.recordWriteOpOverflow(); - } // Execute individually if we're failing to enqueue internalBatchOperation(Collections.singletonList(op)); return; @@ -187,10 +182,5 @@ private void internalBatchOperation(List ops) { this.batchOperation(ops); } - private enum OpType { - READ, - WRITE - } - protected abstract void batchOperation(List ops); } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index dc74a640b22a1..8ec0f43c4914c 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -18,34 +18,20 @@ */ package org.apache.pulsar.metadata.impl.stats; -import io.prometheus.client.Counter; import io.prometheus.client.Gauge; import io.prometheus.client.Histogram; import java.util.concurrent.ExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.atomic.AtomicBoolean; -import org.apache.pulsar.metadata.impl.batching.MetadataOp; -import org.jctools.queues.MessagePassingQueue; public final class BatchMetadataStoreStats implements AutoCloseable { private static final double[] BUCKETS = new double[]{1, 5, 10, 20, 50, 100, 200, 500, 1000}; private static final String NAME = "name"; - private static final String TYPE = "type"; - private static final String OP_READ = "read"; - private static final String OP_WRITE = "write"; - private static final Gauge QUEUEING_OPS = Gauge - .build("pulsar_batch_metadata_store_queueing_ops", "-") - .labelNames(NAME, TYPE) - .register(); private static final Gauge EXECUTOR_QUEUE_SIZE = Gauge .build("pulsar_batch_metadata_store_executor_queue_size", "-") .labelNames(NAME) .register(); - private static final Counter OVERFLOW_OPS = Counter - .build("pulsar_batch_metadata_store_overflow_ops" , "-") - .labelNames(NAME, TYPE) - .register(); private static final Histogram BATCH_OPS_WAITING = Histogram .build("pulsar_batch_metadata_store_queue_wait_time", "-") .unit("ms") @@ -55,40 +41,18 @@ public final class BatchMetadataStoreStats implements AutoCloseable { private final AtomicBoolean closed = new AtomicBoolean(false); private final ThreadPoolExecutor executor; - private final MessagePassingQueue readOps; - private final MessagePassingQueue writeOps; private final String metadataStoreName; - private final Counter.Child readOpsOverflowChild; - private final Counter.Child writeOpsOverflowChild; private final Histogram.Child batchOpsWaitingChild; - public BatchMetadataStoreStats(String metadataStoreName, ExecutorService executor, - MessagePassingQueue readOps, MessagePassingQueue writeOps) { + public BatchMetadataStoreStats(String metadataStoreName, ExecutorService executor) { if (executor instanceof ThreadPoolExecutor tx) { this.executor = tx; } else { this.executor = null; } - this.readOps = readOps; - this.writeOps = writeOps; this.metadataStoreName = metadataStoreName; - - QUEUEING_OPS.setChild(new Gauge.Child() { - @Override - public double get() { - return BatchMetadataStoreStats.this.readOps.size(); - } - }, metadataStoreName, OP_READ); - - QUEUEING_OPS.setChild(new Gauge.Child() { - @Override - public double get() { - return BatchMetadataStoreStats.this.writeOps.size(); - } - }, metadataStoreName, OP_WRITE); - EXECUTOR_QUEUE_SIZE.setChild(new Gauge.Child() { @Override public double get() { @@ -97,8 +61,6 @@ public double get() { } }, metadataStoreName); - this.readOpsOverflowChild = OVERFLOW_OPS.labels(metadataStoreName, OP_READ); - this.writeOpsOverflowChild = OVERFLOW_OPS.labels(metadataStoreName, OP_WRITE); this.batchOpsWaitingChild = BATCH_OPS_WAITING.labels(metadataStoreName); } @@ -106,23 +68,11 @@ public void recordOpWaiting(long millis) { this.batchOpsWaitingChild.observe(millis); } - public void recordReadOpOverflow() { - this.readOpsOverflowChild.inc(); - } - - public void recordWriteOpOverflow() { - this.writeOpsOverflowChild.inc(); - } - @Override public void close() throws Exception { if (closed.compareAndSet(false, true)) { - QUEUEING_OPS.remove(this.metadataStoreName, OP_READ); - QUEUEING_OPS.remove(this.metadataStoreName, OP_WRITE); EXECUTOR_QUEUE_SIZE.remove(this.metadataStoreName); BATCH_OPS_WAITING.remove(this.metadataStoreName); - OVERFLOW_OPS.remove(this.metadataStoreName, OP_READ); - OVERFLOW_OPS.remove(this.metadataStoreName, OP_WRITE); } } } From 2eba39d2ed58a1fd8be16c31e5742f6615922ec0 Mon Sep 17 00:00:00 2001 From: daojun Date: Sat, 29 Oct 2022 21:43:49 +0800 Subject: [PATCH 07/10] fix --- .../broker/stats/MetadataStoreStatsTest.java | 4 +++ .../AbstractBatchedMetadataStore.java | 2 ++ .../impl/stats/BatchMetadataStoreStats.java | 30 +++++++++++++++++-- 3 files changed, 33 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java index 4b47be672d6a9..3a2dcbfecf596 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java @@ -171,9 +171,13 @@ public void testBatchMetadataStoreMetrics() throws Exception { Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_queue_wait_time_ms" + "_sum"); + Collection batchExecuteTime = metricsMap.get("pulsar_batch_metadata_store_batch_execute_time_ms" + "_sum"); + Collection opsPerBatch = metricsMap.get("pulsar_batch_metadata_store_perbatch_ops" + "_sum"); Assert.assertTrue(executorQueueSize.size() > 1); Assert.assertTrue(opsWaiting.size() > 1); + Assert.assertTrue(batchExecuteTime.size() > 0); + Assert.assertTrue(opsPerBatch.size() > 0); for (PrometheusMetricsTest.Metric m : executorQueueSize) { Assert.assertEquals(m.tags.get("cluster"), "test"); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java index 118b553bac985..109ebf82cb198 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/batching/AbstractBatchedMetadataStore.java @@ -180,6 +180,8 @@ private void internalBatchOperation(List ops) { this.batchMetadataStoreStats.recordOpWaiting(now - op.created()); } this.batchOperation(ops); + this.batchMetadataStoreStats.recordOpsInBatch(ops.size()); + this.batchMetadataStoreStats.recordBatchExecuteTime(System.currentTimeMillis() - now); } protected abstract void batchOperation(List ops); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index 8ec0f43c4914c..a65855fff0349 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -32,18 +32,31 @@ public final class BatchMetadataStoreStats implements AutoCloseable { .build("pulsar_batch_metadata_store_executor_queue_size", "-") .labelNames(NAME) .register(); - private static final Histogram BATCH_OPS_WAITING = Histogram + private static final Histogram OPS_WAITING = Histogram .build("pulsar_batch_metadata_store_queue_wait_time", "-") .unit("ms") .labelNames(NAME) .buckets(BUCKETS) .register(); + private static final Histogram BATCH_EXECUTE_TIME = Histogram + .build("pulsar_batch_metadata_store_batch_execute_time", "-") + .unit("ms") + .labelNames(NAME) + .buckets(BUCKETS) + .register(); + private static final Histogram OPS_PER_BATCH = Histogram + .build("pulsar_batch_metadata_store_perbatch_ops", "-") + .labelNames(NAME) + .buckets(BUCKETS) + .register(); private final AtomicBoolean closed = new AtomicBoolean(false); private final ThreadPoolExecutor executor; private final String metadataStoreName; private final Histogram.Child batchOpsWaitingChild; + private final Histogram.Child batchExecuteTimeChild; + private final Histogram.Child opsPerBatchChild; public BatchMetadataStoreStats(String metadataStoreName, ExecutorService executor) { if (executor instanceof ThreadPoolExecutor tx) { @@ -61,18 +74,29 @@ public double get() { } }, metadataStoreName); - this.batchOpsWaitingChild = BATCH_OPS_WAITING.labels(metadataStoreName); + this.batchOpsWaitingChild = OPS_WAITING.labels(metadataStoreName); + this.batchExecuteTimeChild = BATCH_EXECUTE_TIME.labels(metadataStoreName); + this.opsPerBatchChild = OPS_PER_BATCH.labels(metadataStoreName); + } public void recordOpWaiting(long millis) { this.batchOpsWaitingChild.observe(millis); } + public void recordBatchExecuteTime(long millis) { + this.batchExecuteTimeChild.observe(millis); + } + + public void recordOpsInBatch(int ops) { + this.opsPerBatchChild.observe(ops); + } + @Override public void close() throws Exception { if (closed.compareAndSet(false, true)) { EXECUTOR_QUEUE_SIZE.remove(this.metadataStoreName); - BATCH_OPS_WAITING.remove(this.metadataStoreName); + OPS_WAITING.remove(this.metadataStoreName); } } } From 8f8719c6364b127b65a3cab1940b53585d3dae6c Mon Sep 17 00:00:00 2001 From: daojun Date: Sat, 29 Oct 2022 22:06:14 +0800 Subject: [PATCH 08/10] fix tests --- .../broker/stats/MetadataStoreStatsTest.java | 20 +++++++++++++++++++ .../impl/stats/BatchMetadataStoreStats.java | 2 +- 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java index 00070cb9dd665..4a41e51707bfe 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java @@ -197,6 +197,26 @@ public void testBatchMetadataStoreMetrics() throws Exception { || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); Assert.assertTrue(m.value >= 0); } + + for (PrometheusMetricsTest.Metric m : batchExecuteTime) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + String metadataStoreName = m.tags.get("name"); + Assert.assertNotNull(metadataStoreName); + Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); + Assert.assertTrue(m.value > 0); + } + + for (PrometheusMetricsTest.Metric m : opsPerBatch) { + Assert.assertEquals(m.tags.get("cluster"), "test"); + String metadataStoreName = m.tags.get("name"); + Assert.assertNotNull(metadataStoreName); + Assert.assertTrue(metadataStoreName.equals(MetadataStoreConfig.METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.CONFIGURATION_METADATA_STORE) + || metadataStoreName.equals(MetadataStoreConfig.STATE_METADATA_STORE)); + Assert.assertTrue(m.value > 0); + } } } diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index a65855fff0349..f76e782b72e0e 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -1,4 +1,4 @@ -/** +/* * 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 From 316bcd0a99cf815bcc9eed1d022e191c3b095472 Mon Sep 17 00:00:00 2001 From: daojun Date: Tue, 1 Nov 2022 01:52:54 +0800 Subject: [PATCH 09/10] fix metrics name --- .../org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java | 2 +- .../pulsar/metadata/impl/stats/BatchMetadataStoreStats.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java index 4a41e51707bfe..3ce735b797fb1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/MetadataStoreStatsTest.java @@ -172,7 +172,7 @@ public void testBatchMetadataStoreMetrics() throws Exception { Collection executorQueueSize = metricsMap.get("pulsar_batch_metadata_store_executor_queue_size"); Collection opsWaiting = metricsMap.get("pulsar_batch_metadata_store_queue_wait_time_ms" + "_sum"); Collection batchExecuteTime = metricsMap.get("pulsar_batch_metadata_store_batch_execute_time_ms" + "_sum"); - Collection opsPerBatch = metricsMap.get("pulsar_batch_metadata_store_perbatch_ops" + "_sum"); + Collection opsPerBatch = metricsMap.get("pulsar_batch_metadata_store_batch_size" + "_sum"); Assert.assertTrue(executorQueueSize.size() > 1); Assert.assertTrue(opsWaiting.size() > 1); diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index f76e782b72e0e..ed3750777840e 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -45,7 +45,7 @@ public final class BatchMetadataStoreStats implements AutoCloseable { .buckets(BUCKETS) .register(); private static final Histogram OPS_PER_BATCH = Histogram - .build("pulsar_batch_metadata_store_perbatch_ops", "-") + .build("pulsar_batch_metadata_store_batch_size", "-") .labelNames(NAME) .buckets(BUCKETS) .register(); From 5127d2523fafe1dd9083a7d9ca0e14b54b0f793f Mon Sep 17 00:00:00 2001 From: daojun Date: Tue, 1 Nov 2022 15:46:04 +0800 Subject: [PATCH 10/10] fix metrics name --- .../pulsar/metadata/impl/stats/BatchMetadataStoreStats.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java index ed3750777840e..f87155b9259be 100644 --- a/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java +++ b/pulsar-metadata/src/main/java/org/apache/pulsar/metadata/impl/stats/BatchMetadataStoreStats.java @@ -97,6 +97,8 @@ public void close() throws Exception { if (closed.compareAndSet(false, true)) { EXECUTOR_QUEUE_SIZE.remove(this.metadataStoreName); OPS_WAITING.remove(this.metadataStoreName); + BATCH_EXECUTE_TIME.remove(this.metadataStoreName); + OPS_PER_BATCH.remove(metadataStoreName); } } }