diff --git a/conf/broker.conf b/conf/broker.conf index 76f75a7d82530..ec178720113e7 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -1461,9 +1461,8 @@ exposePreciseBacklogInPrometheus=false splitTopicAndPartitionLabelInPrometheus=false -# If true and the client supports partial producer, aggregate publisher stats of PartitionedTopicStats by producerName. -# Otherwise, aggregate it by list index. -aggregatePublisherStatsByProducerName=false +# Deprecated: it aggregates publisher stats by producerName +aggregatePublisherStatsByProducerName= # Interval between checks to see if cluster is migrated and marks topic migrated # if cluster is marked migrated. Disable with value 0. (Default disabled). diff --git a/conf/standalone.conf b/conf/standalone.conf index b1b3a068e35c2..13c7ea124cdbf 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -987,9 +987,9 @@ exposePreciseBacklogInPrometheus=false splitTopicAndPartitionLabelInPrometheus=false -# If true, aggregate publisher stats of PartitionedTopicStats by producerName. -# Otherwise, aggregate it by list index. -aggregatePublisherStatsByProducerName=false + +# Deprecated: it aggregates publisher stats by producerName +aggregatePublisherStatsByProducerName= ### --- Schema storage --- ### # The schema storage implementation used by this broker. 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 f0725353dd717..38f8e37a9f296 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 @@ -2705,11 +2705,6 @@ The delayed message index bucket time step(in seconds) in per bucket snapshot se doc = "Stats update initial delay in seconds" ) private int statsUpdateInitialDelayInSecs = 60; - @FieldContext( - category = CATEGORY_METRICS, - doc = "If true, aggregate publisher stats of PartitionedTopicStats by producerName" - ) - private boolean aggregatePublisherStatsByProducerName = false; /**** --- Ledger Offloading. --- ****/ /**** diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index d7b99f9e8146e..67e00fb36afcd 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -100,8 +100,7 @@ public Producer(Topic topic, TransportCnx cnx, long producerId, String producerN boolean isEncrypted, Map metadata, SchemaVersion schemaVersion, long epoch, boolean userProvidedProducerName, ProducerAccessMode accessMode, - Optional topicEpoch, - boolean supportsPartialProducer) { + Optional topicEpoch) { final ServiceConfiguration serviceConf = cnx.getBrokerService().pulsar().getConfiguration(); this.topic = topic; @@ -129,15 +128,6 @@ public Producer(Topic topic, TransportCnx cnx, long producerId, String producerN stats.setClientVersion(cnx.getClientVersion()); stats.setProducerName(producerName); stats.producerId = producerId; - if (serviceConf.isAggregatePublisherStatsByProducerName() && stats.getProducerName() != null) { - // If true and the client supports partial producer, - // aggregate publisher stats of PartitionedTopicStats by producerName. - // Otherwise, aggregate it by list index. - stats.setSupportsPartialProducer(supportsPartialProducer); - } else { - // aggregate publisher stats of PartitionedTopicStats by list index. - stats.setSupportsPartialProducer(false); - } stats.metadata = this.metadata; stats.accessMode = Commands.convertProducerAccessMode(accessMode); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 4c71ecee2e288..b14d71c50115d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1236,7 +1236,6 @@ protected void handleProducer(final CommandProducer cmdProducer) { final boolean isTxnEnabled = cmdProducer.isTxnEnabled(); final String initialSubscriptionName = cmdProducer.hasInitialSubscriptionName() ? cmdProducer.getInitialSubscriptionName() : null; - final boolean supportsPartialProducer = supportsPartialProducer(); TopicName topicName = validateTopicName(cmdProducer.getTopic(), requestId, cmdProducer); if (topicName == null) { @@ -1392,7 +1391,7 @@ protected void handleProducer(final CommandProducer cmdProducer) { buildProducerAndAddTopic(topic, producerId, producerName, requestId, isEncrypted, metadata, schemaVersion, epoch, userProvidedProducerName, topicName, - producerAccessMode, topicEpoch, supportsPartialProducer, producerFuture); + producerAccessMode, topicEpoch, producerFuture); }); }).exceptionally(exception -> { Throwable cause = exception.getCause(); @@ -1464,12 +1463,12 @@ private void buildProducerAndAddTopic(Topic topic, long producerId, String produ boolean isEncrypted, Map metadata, SchemaVersion schemaVersion, long epoch, boolean userProvidedProducerName, TopicName topicName, ProducerAccessMode producerAccessMode, - Optional topicEpoch, boolean supportsPartialProducer, + Optional topicEpoch, CompletableFuture producerFuture){ CompletableFuture producerQueuedFuture = new CompletableFuture<>(); Producer producer = new Producer(topic, ServerCnx.this, producerId, producerName, getPrincipal(), isEncrypted, metadata, schemaVersion, epoch, - userProvidedProducerName, producerAccessMode, topicEpoch, supportsPartialProducer); + userProvidedProducerName, producerAccessMode, topicEpoch); topic.addProducer(producer, producerQueuedFuture).thenAccept(newTopicEpoch -> { if (isActive()) { @@ -3016,10 +3015,6 @@ boolean supportBrokerMetadata() { return features != null && features.isSupportsBrokerEntryMetadata(); } - boolean supportsPartialProducer() { - return features != null && features.isSupportsPartialProducer(); - } - @Override public String getClientVersion() { return clientVersion; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index d04d0c12d6f69..c919d279954f5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -2696,7 +2696,6 @@ public void testPartitionedStatsAggregationByProducerName(String topicType) thro cleanup(); setup(); - conf.setAggregatePublisherStatsByProducerName(true); final String topic = topicType + "://prop-xyz/ns1/test-partitioned-stats-aggregation-by-producer-name"; admin.topics().createPartitionedTopic(topic, 10); @@ -2746,12 +2745,10 @@ public int choosePartition(Message msg, TopicMetadata metadata) { assertEquals(topicStats.getPartitions().values().stream().mapToInt(e -> e.getPublishers().size()).sum(), 10); assertEquals(topicStats.getPartitions().values().stream().map(e -> e.getPublishers().get(0).getProducerName()).distinct().count(), 2); assertEquals(topicStats.getPublishers().size(), 2); - topicStats.getPublishers().forEach(p -> assertTrue(p.isSupportsPartialProducer())); } @Test(dataProvider = "topicType") public void testPartitionedStatsAggregationByProducerNamePerPartition(String topicType) throws Exception { - conf.setAggregatePublisherStatsByProducerName(true); final String topic = topicType + "://prop-xyz/ns1/test-partitioned-stats-aggregation-by-producer-name-per-pt"; admin.topics().createPartitionedTopic(topic, 2); @@ -2770,7 +2767,6 @@ public void testPartitionedStatsAggregationByProducerNamePerPartition(String top assertEquals(topicStats.getPartitions().values().stream().mapToInt(e -> e.getPublishers().size()).sum(), 2); assertEquals(topicStats.getPartitions().values().stream().map(e -> e.getPublishers().get(0).getProducerName()).distinct().count(), 2); assertEquals(topicStats.getPublishers().size(), 2); - topicStats.getPublishers().forEach(p -> assertTrue(p.isSupportsPartialProducer())); } @Test(dataProvider = "topicType") diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index 41f3dee76468d..669c1e79d87ae 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -438,7 +438,7 @@ public void testAddRemoveProducer() throws Exception { // 1. simple add producer Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", role, false, null, SchemaVersion.Latest, 0, false, - ProducerAccessMode.Shared, Optional.empty(), true); + ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer, new CompletableFuture<>()); assertEquals(topic.getProducers().size(), 1); @@ -455,7 +455,7 @@ public void testAddRemoveProducer() throws Exception { PersistentTopic failTopic = new PersistentTopic(failTopicName, ledgerMock, brokerService); Producer failProducer = new Producer(failTopic, serverCnx, 2 /* producer id */, "prod-name", role, false, null, SchemaVersion.Latest, 0, false, - ProducerAccessMode.Shared, Optional.empty(), true); + ProducerAccessMode.Shared, Optional.empty()); try { topic.addProducer(failProducer, new CompletableFuture<>()); fail("should have failed"); @@ -466,7 +466,7 @@ public void testAddRemoveProducer() throws Exception { // 4. Try to remove with unequal producer Producer producerCopy = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", role, false, null, SchemaVersion.Latest, 0, false, - ProducerAccessMode.Shared, Optional.empty(), true); + ProducerAccessMode.Shared, Optional.empty()); topic.removeProducer(producerCopy); // Expect producer to be in map assertEquals(topic.getProducers().size(), 1); @@ -485,9 +485,9 @@ public void testProducerOverwrite() throws Exception { PersistentTopic topic = new PersistentTopic(successTopicName, ledgerMock, brokerService); String role = "appid1"; Producer producer1 = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, true, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 0, true, ProducerAccessMode.Shared, Optional.empty()); Producer producer2 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, true, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 0, true, ProducerAccessMode.Shared, Optional.empty()); try { topic.addProducer(producer1, new CompletableFuture<>()).join(); topic.addProducer(producer2, new CompletableFuture<>()).join(); @@ -500,7 +500,7 @@ public void testProducerOverwrite() throws Exception { Assert.assertEquals(topic.getProducers().size(), 1); Producer producer3 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 1, false, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 1, false, ProducerAccessMode.Shared, Optional.empty()); try { topic.addProducer(producer3, new CompletableFuture<>()).join(); @@ -516,7 +516,7 @@ public void testProducerOverwrite() throws Exception { Assert.assertEquals(topic.getProducers().size(), 0); Producer producer4 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 2, false, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 2, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer3, new CompletableFuture<>()); topic.addProducer(producer4, new CompletableFuture<>()); @@ -529,13 +529,13 @@ public void testProducerOverwrite() throws Exception { Assert.assertEquals(topic.getProducers().size(), 0); Producer producer5 = new Producer(topic, serverCnx, 2 /* producer id */, "pulsar.repl.cluster1", - role, false, null, SchemaVersion.Latest, 1, false, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 1, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer5, new CompletableFuture<>()); Assert.assertEquals(topic.getProducers().size(), 1); Producer producer6 = new Producer(topic, serverCnx, 2 /* producer id */, "pulsar.repl.cluster1", - role, false, null, SchemaVersion.Latest, 2, false, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 2, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer6, new CompletableFuture<>()); Assert.assertEquals(topic.getProducers().size(), 1); @@ -543,7 +543,7 @@ public void testProducerOverwrite() throws Exception { topic.getProducers().values().forEach(producer -> Assert.assertEquals(producer.getEpoch(), 2)); Producer producer7 = new Producer(topic, serverCnx, 2 /* producer id */, "pulsar.repl.cluster1", - role, false, null, SchemaVersion.Latest, 3, true, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 3, true, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer7, new CompletableFuture<>()); Assert.assertEquals(topic.getProducers().size(), 1); @@ -556,20 +556,20 @@ private void testMaxProducers() throws Exception { String role = "appid1"; // 1. add producer1 Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name1", role, - false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty(), true); + false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer, new CompletableFuture<>()); assertEquals(topic.getProducers().size(), 1); // 2. add producer2 Producer producer2 = new Producer(topic, serverCnx, 2 /* producer id */, "prod-name2", role, - false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty(), true); + false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer2, new CompletableFuture<>()); assertEquals(topic.getProducers().size(), 2); // 3. add producer3 but reached maxProducersPerTopic try { Producer producer3 = new Producer(topic, serverCnx, 3 /* producer id */, "prod-name3", role, - false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty(), true); + false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer3, new CompletableFuture<>()).join(); fail("should have failed"); } catch (Exception e) { @@ -619,7 +619,7 @@ private Producer getMockedProducerWithSpecificAddress(Topic topic, long producer doReturn(new PulsarCommandSenderImpl(null, cnx)).when(cnx).getCommandSender(); return new Producer(topic, cnx, producerId, producerNameBase + producerId, role, false, null, - SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty(), true); + SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty()); } @Test @@ -1292,7 +1292,7 @@ public void testDeleteTopic() throws Exception { // 3. delete topic with producer topic = (PersistentTopic) brokerService.getOrCreateTopic(successTopicName).get(); Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer, new CompletableFuture<>()).join(); assertTrue(topic.delete().isCompletedExceptionally()); @@ -1466,7 +1466,7 @@ public Object answer(InvocationOnMock invocationOnMock) throws Throwable { String role = "appid1"; Thread.sleep(10); /* delay to ensure that the delete gets executed first */ Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", - role, false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty(), true); + role, false, null, SchemaVersion.Latest, 0, false, ProducerAccessMode.Shared, Optional.empty()); topic.addProducer(producer, new CompletableFuture<>()).join(); fail("Should have failed"); } catch (Exception e) { @@ -2253,7 +2253,7 @@ public void testDisconnectProducer() throws Exception { String role = "appid1"; Producer producer = new Producer(topic, serverCnx, 1 /* producer id */, "prod-name", role, false, null, SchemaVersion.Latest, 0, false, - ProducerAccessMode.Shared, Optional.empty(), true); + ProducerAccessMode.Shared, Optional.empty()); assertFalse(producer.isDisconnecting()); // Disconnect the producer multiple times. producer.disconnect(); diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/PublisherStats.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/PublisherStats.java index 78af224d28c02..af7d5130e4d5f 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/PublisherStats.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/PublisherStats.java @@ -44,6 +44,7 @@ public interface PublisherStats { long getProducerId(); /** Whether partial producer is supported at client. */ + @Deprecated boolean isSupportsPartialProducer(); /** Producer name. */ diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/NonPersistentTopicStatsImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/NonPersistentTopicStatsImpl.java index 06e955ef3d6a7..386182608a01a 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/NonPersistentTopicStatsImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/NonPersistentTopicStatsImpl.java @@ -18,13 +18,9 @@ */ package org.apache.pulsar.common.policies.data.stats; -import static java.util.Comparator.naturalOrder; -import static java.util.Comparator.nullsLast; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; -import java.util.ArrayList; -import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -32,13 +28,11 @@ import java.util.TreeMap; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; -import java.util.stream.Stream; import lombok.Getter; import org.apache.pulsar.common.policies.data.NonPersistentPublisherStats; import org.apache.pulsar.common.policies.data.NonPersistentReplicatorStats; import org.apache.pulsar.common.policies.data.NonPersistentSubscriptionStats; import org.apache.pulsar.common.policies.data.NonPersistentTopicStats; -import org.apache.pulsar.common.policies.data.PublisherStats; /** * Statistics for a non-persistent topic. @@ -53,9 +47,6 @@ public class NonPersistentTopicStatsImpl extends TopicStatsImpl implements NonPe @Getter public double msgDropRate; - @JsonIgnore - public List publishers; - @JsonIgnore public Map subscriptions; @@ -64,11 +55,7 @@ public class NonPersistentTopicStatsImpl extends TopicStatsImpl implements NonPe @JsonProperty("publishers") public List getNonPersistentPublishers() { - return Stream.concat(nonPersistentPublishers.stream().sorted( - Comparator.comparing(NonPersistentPublisherStats::getProducerName, nullsLast(naturalOrder()))), - nonPersistentPublishersMap.values().stream().sorted( - Comparator.comparing(NonPersistentPublisherStats::getProducerName, nullsLast(naturalOrder())))) - .collect(Collectors.toList()); + return nonPersistentPublishers.values().stream().collect(Collectors.toList()); } @JsonProperty("subscriptions") @@ -81,10 +68,8 @@ public Map getNonPersistentReplicators() { return (Map) nonPersistentReplicators; } - /** List of connected publishers on this non-persistent topic w/ their stats. */ - private List nonPersistentPublishers; - - private Map nonPersistentPublishersMap; + /** Map of connected publishers on this non-persistent topic w/ their stats. */ + private Map nonPersistentPublishers; /** Map of non-persistent subscriptions with their individual statistics. */ public Map nonPersistentSubscriptions; @@ -94,26 +79,17 @@ public Map getNonPersistentReplicators() { @SuppressFBWarnings(value = "MF_CLASS_MASKS_FIELD", justification = "expected to override") public List getPublishers() { - return Stream.concat(nonPersistentPublishers.stream().sorted( - Comparator.comparing(NonPersistentPublisherStats::getProducerName, nullsLast(naturalOrder()))), - nonPersistentPublishersMap.values().stream().sorted( - Comparator.comparing(NonPersistentPublisherStats::getProducerName, nullsLast(naturalOrder())))) - .collect(Collectors.toList()); + return nonPersistentPublishers.values().stream().collect(Collectors.toList()); } - public void setPublishers(List statsList) { + @JsonProperty("publishers") + public void setNonPersistentPublishers(List statsList) { this.nonPersistentPublishers.clear(); - this.nonPersistentPublishersMap.clear(); - statsList.forEach(s -> addPublisher((NonPersistentPublisherStatsImpl) s)); + statsList.forEach(s -> addPublisher(s)); } - public void addPublisher(NonPersistentPublisherStatsImpl stats) { - if (stats.isSupportsPartialProducer() && stats.getProducerName() != null) { - nonPersistentPublishersMap.put(stats.getProducerName(), stats); - } else { - stats.setSupportsPartialProducer(false); // setter method with side effect - nonPersistentPublishers.add(stats); - } + public void addPublisher(NonPersistentPublisherStats stats) { + nonPersistentPublishers.put(stats.getProducerName(), stats); } @SuppressFBWarnings(value = "MF_CLASS_MASKS_FIELD", justification = "expected to override") @@ -132,8 +108,7 @@ public double getMsgDropRate() { } public NonPersistentTopicStatsImpl() { - this.nonPersistentPublishers = new ArrayList<>(); - this.nonPersistentPublishersMap = new ConcurrentHashMap<>(); + this.nonPersistentPublishers = new ConcurrentHashMap<>(); this.nonPersistentSubscriptions = new HashMap<>(); this.nonPersistentReplicators = new TreeMap<>(); } @@ -141,7 +116,6 @@ public NonPersistentTopicStatsImpl() { public void reset() { super.reset(); this.nonPersistentPublishers.clear(); - this.nonPersistentPublishersMap.clear(); this.nonPersistentSubscriptions.clear(); this.nonPersistentReplicators.clear(); this.msgDropRate = 0; @@ -155,71 +129,30 @@ public NonPersistentTopicStatsImpl add(NonPersistentTopicStats ts) { super.add(stats); this.msgDropRate += stats.msgDropRate; - stats.getNonPersistentPublishers().forEach(s -> { - if (s.isSupportsPartialProducer() && s.getProducerName() != null) { - ((NonPersistentPublisherStatsImpl) this.nonPersistentPublishersMap - .computeIfAbsent(s.getProducerName(), key -> { + // Aggregate the input publish stats to this.nonPersistentPublishers grouped by producer name + stats.nonPersistentPublishers.values().forEach(pubStat -> + ((NonPersistentPublisherStatsImpl) this.nonPersistentPublishers + .computeIfAbsent(pubStat.getProducerName(), key -> { final NonPersistentPublisherStatsImpl newStats = new NonPersistentPublisherStatsImpl(); - newStats.setSupportsPartialProducer(true); - newStats.setProducerName(s.getProducerName()); + newStats.setProducerName(pubStat.getProducerName()); return newStats; - })).add((NonPersistentPublisherStatsImpl) s); - } else { - if (this.nonPersistentPublishers.size() != stats.getNonPersistentPublishers().size()) { - for (int i = 0; i < stats.getNonPersistentPublishers().size(); i++) { - NonPersistentPublisherStatsImpl newStats = new NonPersistentPublisherStatsImpl(); - newStats.setSupportsPartialProducer(false); - this.nonPersistentPublishers.add(newStats.add((NonPersistentPublisherStatsImpl) s)); - } - } else { - for (int i = 0; i < stats.getNonPersistentPublishers().size(); i++) { - ((NonPersistentPublisherStatsImpl) this.nonPersistentPublishers.get(i)) - .add((NonPersistentPublisherStatsImpl) s); - } - } - } - }); - - if (this.getNonPersistentSubscriptions().size() != stats.getNonPersistentSubscriptions().size()) { - for (String subscription : stats.getNonPersistentSubscriptions().keySet()) { - NonPersistentSubscriptionStatsImpl subscriptionStats = new NonPersistentSubscriptionStatsImpl(); - this.getNonPersistentSubscriptions().put(subscription, subscriptionStats - .add((NonPersistentSubscriptionStatsImpl) - stats.getNonPersistentSubscriptions().get(subscription))); - } - } else { - for (String subscription : stats.getNonPersistentSubscriptions().keySet()) { - if (this.getNonPersistentSubscriptions().get(subscription) != null) { - ((NonPersistentSubscriptionStatsImpl) this.getNonPersistentSubscriptions().get(subscription)) - .add((NonPersistentSubscriptionStatsImpl) - stats.getNonPersistentSubscriptions().get(subscription)); - } else { - NonPersistentSubscriptionStatsImpl subscriptionStats = new NonPersistentSubscriptionStatsImpl(); - this.getNonPersistentSubscriptions().put(subscription, subscriptionStats - .add((NonPersistentSubscriptionStatsImpl) - stats.getNonPersistentSubscriptions().get(subscription))); - } - } - } - - if (this.getNonPersistentReplicators().size() != stats.getNonPersistentReplicators().size()) { - for (String repl : stats.getNonPersistentReplicators().keySet()) { - NonPersistentReplicatorStatsImpl replStats = new NonPersistentReplicatorStatsImpl(); - this.getNonPersistentReplicators().put(repl, replStats - .add((NonPersistentReplicatorStatsImpl) stats.getNonPersistentReplicators().get(repl))); - } - } else { - for (String repl : stats.getNonPersistentReplicators().keySet()) { - if (this.getNonPersistentReplicators().get(repl) != null) { - ((NonPersistentReplicatorStatsImpl) this.getNonPersistentReplicators().get(repl)) - .add((NonPersistentReplicatorStatsImpl) stats.getNonPersistentReplicators().get(repl)); - } else { - NonPersistentReplicatorStatsImpl replStats = new NonPersistentReplicatorStatsImpl(); - this.getNonPersistentReplicators().put(repl, replStats - .add((NonPersistentReplicatorStatsImpl) stats.getNonPersistentReplicators().get(repl))); - } - } - } + })) + .add((NonPersistentPublisherStatsImpl) pubStat) + ); + + // Aggregate the input subscription stats to this.nonPersistentSubscriptions grouped by subscription name + stats.getNonPersistentSubscriptions().entrySet().forEach(subscription -> + ((NonPersistentSubscriptionStatsImpl) this.getNonPersistentSubscriptions() + .computeIfAbsent(subscription.getKey(), key -> new NonPersistentSubscriptionStatsImpl())) + .add((NonPersistentSubscriptionStatsImpl) subscription.getValue()) + ); + + // Aggregate the input replicator stats to this.nonPersistentReplicators grouped by replicator name + stats.getNonPersistentReplicators().entrySet().forEach(replicator -> + ((NonPersistentReplicatorStatsImpl) this.getNonPersistentReplicators() + .computeIfAbsent(replicator.getKey(), key -> new NonPersistentReplicatorStatsImpl())) + .add((NonPersistentReplicatorStatsImpl) replicator.getValue()) + ); return this; } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/PublisherStatsImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/PublisherStatsImpl.java index 41407a37e7ca0..e29ded08d7faa 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/PublisherStatsImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/PublisherStatsImpl.java @@ -49,8 +49,9 @@ public class PublisherStatsImpl implements PublisherStats { /** Id of this publisher. */ public long producerId; - /** Whether partial producer is supported at client. */ - public boolean supportsPartialProducer; + /** Whether partial producer is supported at client. supportsPartialProducer is true always. */ + @Deprecated + public boolean supportsPartialProducer = true; /** Producer name. */ @JsonIgnore @@ -99,7 +100,7 @@ public PublisherStatsImpl add(PublisherStatsImpl stats) { } public String getProducerName() { - return producerNameOffset == -1 ? null + return producerNameOffset == -1 ? "null" : stringBuffer.substring(producerNameOffset, producerNameOffset + producerNameLength); } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/TopicStatsImpl.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/TopicStatsImpl.java index 88f487f347a9a..4b9eeda76d567 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/TopicStatsImpl.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/stats/TopicStatsImpl.java @@ -18,18 +18,13 @@ */ package org.apache.pulsar.common.policies.data.stats; -import static java.util.Comparator.naturalOrder; -import static java.util.Comparator.nullsLast; import com.fasterxml.jackson.annotation.JsonIgnore; -import java.util.ArrayList; -import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.TreeMap; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; -import java.util.stream.Stream; import lombok.AccessLevel; import lombok.Data; import lombok.Getter; @@ -105,14 +100,10 @@ public class TopicStatsImpl implements TopicStats { public long abortedTxnCount; public long committedTxnCount; - /** List of connected publishers on this topic w/ their stats. */ + /** Map of connected publishers on this topic w/ their stats. */ @Getter(AccessLevel.NONE) @Setter(AccessLevel.NONE) - private List publishers; - - @Getter(AccessLevel.NONE) - @Setter(AccessLevel.NONE) - private Map publishersMap; + private Map publishers; public int waitingPublishers; @@ -145,26 +136,16 @@ public class TopicStatsImpl implements TopicStats { public String ownerBroker; public List getPublishers() { - return Stream.concat(publishers.stream().sorted( - Comparator.comparing(PublisherStatsImpl::getProducerName, nullsLast(naturalOrder()))), - publishersMap.values().stream().sorted( - Comparator.comparing(PublisherStatsImpl::getProducerName, nullsLast(naturalOrder())))) - .collect(Collectors.toList()); + return publishers.values().stream().collect(Collectors.toList()); } public void setPublishers(List statsList) { this.publishers.clear(); - this.publishersMap.clear(); statsList.forEach(s -> addPublisher((PublisherStatsImpl) s)); } - public void addPublisher(PublisherStatsImpl stats) { - if (stats.isSupportsPartialProducer() && stats.getProducerName() != null) { - publishersMap.put(stats.getProducerName(), stats); - } else { - stats.setSupportsPartialProducer(false); // setter method with side effect - publishers.add(stats); - } + public void addPublisher(PublisherStats stats) { + publishers.put(stats.getProducerName(), (PublisherStatsImpl) stats); } public Map getSubscriptions() { @@ -176,8 +157,7 @@ public void addPublisher(PublisherStatsImpl stats) { } public TopicStatsImpl() { - this.publishers = new ArrayList<>(); - this.publishersMap = new ConcurrentHashMap<>(); + this.publishers = new ConcurrentHashMap<>(); this.subscriptions = new HashMap<>(); this.replication = new TreeMap<>(); this.compaction = new CompactionStatsImpl(); @@ -197,7 +177,6 @@ public void reset() { this.bytesOutCounter = 0; this.msgOutCounter = 0; this.publishers.clear(); - this.publishersMap.clear(); this.subscriptions.clear(); this.waitingPublishers = 0; this.replication.clear(); @@ -243,59 +222,27 @@ public TopicStatsImpl add(TopicStats ts) { this.abortedTxnCount = stats.abortedTxnCount; this.committedTxnCount = stats.committedTxnCount; - stats.getPublishers().forEach(s -> { - if (s.isSupportsPartialProducer() && s.getProducerName() != null) { - this.publishersMap.computeIfAbsent(s.getProducerName(), key -> { - final PublisherStatsImpl newStats = new PublisherStatsImpl(); - newStats.setSupportsPartialProducer(true); - newStats.setProducerName(s.getProducerName()); - return newStats; - }).add((PublisherStatsImpl) s); - } else { - if (this.publishers.size() != stats.publishers.size()) { - for (int i = 0; i < stats.publishers.size(); i++) { - PublisherStatsImpl newStats = new PublisherStatsImpl(); - newStats.setSupportsPartialProducer(false); - this.publishers.add(newStats.add(stats.publishers.get(i))); - } - } else { - for (int i = 0; i < stats.publishers.size(); i++) { - this.publishers.get(i).add(stats.publishers.get(i)); - } - } - } - }); - - if (this.subscriptions.size() != stats.subscriptions.size()) { - for (String subscription : stats.subscriptions.keySet()) { - SubscriptionStatsImpl subscriptionStats = new SubscriptionStatsImpl(); - this.subscriptions.put(subscription, subscriptionStats.add(stats.subscriptions.get(subscription))); - } - } else { - for (String subscription : stats.subscriptions.keySet()) { - if (this.subscriptions.get(subscription) != null) { - this.subscriptions.get(subscription).add(stats.subscriptions.get(subscription)); - } else { - SubscriptionStatsImpl subscriptionStats = new SubscriptionStatsImpl(); - this.subscriptions.put(subscription, subscriptionStats.add(stats.subscriptions.get(subscription))); - } - } - } - if (this.replication.size() != stats.replication.size()) { - for (String repl : stats.replication.keySet()) { - ReplicatorStatsImpl replStats = new ReplicatorStatsImpl(); - this.replication.put(repl, replStats.add(stats.replication.get(repl))); - } - } else { - for (String repl : stats.replication.keySet()) { - if (this.replication.get(repl) != null) { - this.replication.get(repl).add(stats.replication.get(repl)); - } else { - ReplicatorStatsImpl replStats = new ReplicatorStatsImpl(); - this.replication.put(repl, replStats.add(stats.replication.get(repl))); - } - } - } + // Aggregate the input publish stats to this.nonPersistentPublishers grouped by producer name + stats.publishers.values().forEach(pubStat -> + this.publishers.computeIfAbsent(pubStat.getProducerName(), key -> { + final PublisherStatsImpl newStats = new PublisherStatsImpl(); + newStats.setProducerName(pubStat.getProducerName()); + return newStats; + }).add(pubStat) + ); + + // Aggregate the input subscription stats to this.nonPersistentSubscriptions grouped by subscription name + stats.subscriptions.entrySet().forEach(subscription -> + this.subscriptions.computeIfAbsent(subscription.getKey(), key -> new SubscriptionStatsImpl()) + .add(subscription.getValue()) + ); + + // Aggregate the input replicator stats to this.nonPersistentReplicators grouped by replicator name + stats.replication.entrySet().forEach(replicator -> + this.replication.computeIfAbsent(replicator.getKey(), key -> new ReplicatorStatsImpl()) + .add(replicator.getValue()) + ); + return this; } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index fdb94c177959c..26f02cf319f7f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -187,7 +187,6 @@ public static ByteBuf newConnect(String authMethodName, String authData, String private static void setFeatureFlags(FeatureFlags flags) { flags.setSupportsAuthRefresh(true); flags.setSupportsBrokerEntryMetadata(true); - flags.setSupportsPartialProducer(true); } public static ByteBuf newConnect(String authMethodName, String authData, int protocolVersion, String libVersion, diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/NonPersistentPartitionedTopicStatsTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/NonPersistentPartitionedTopicStatsTest.java index 814df9e712a11..21088b757f8a7 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/NonPersistentPartitionedTopicStatsTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/NonPersistentPartitionedTopicStatsTest.java @@ -18,15 +18,17 @@ */ package org.apache.pulsar.common.policies.data; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.pulsar.common.policies.data.stats.NonPersistentPartitionedTopicStatsImpl; import org.apache.pulsar.common.policies.data.stats.NonPersistentPublisherStatsImpl; import org.apache.pulsar.common.policies.data.stats.NonPersistentReplicatorStatsImpl; import org.apache.pulsar.common.policies.data.stats.NonPersistentSubscriptionStatsImpl; import org.apache.pulsar.common.policies.data.stats.NonPersistentTopicStatsImpl; +import org.apache.pulsar.common.util.ObjectMapperFactory; import org.testng.annotations.Test; import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertFalse; public class NonPersistentPartitionedTopicStatsTest { @@ -62,37 +64,73 @@ public void testPartitionedTopicStats() { public void testPartitionedTopicStatsByNullProducerName() { final NonPersistentTopicStatsImpl topicStats1 = new NonPersistentTopicStatsImpl(); final NonPersistentPublisherStatsImpl publisherStats1 = new NonPersistentPublisherStatsImpl(); - publisherStats1.setSupportsPartialProducer(false); publisherStats1.setProducerName(null); final NonPersistentPublisherStatsImpl publisherStats2 = new NonPersistentPublisherStatsImpl(); - publisherStats2.setSupportsPartialProducer(false); publisherStats2.setProducerName(null); topicStats1.addPublisher(publisherStats1); topicStats1.addPublisher(publisherStats2); - assertEquals(topicStats1.getPublishers().size(), 2); - assertFalse(topicStats1.getPublishers().get(0).isSupportsPartialProducer()); - assertFalse(topicStats1.getPublishers().get(1).isSupportsPartialProducer()); + assertEquals(topicStats1.getPublishers().size(), 1); final NonPersistentTopicStatsImpl topicStats2 = new NonPersistentTopicStatsImpl(); final NonPersistentPublisherStatsImpl publisherStats3 = new NonPersistentPublisherStatsImpl(); - publisherStats3.setSupportsPartialProducer(true); publisherStats3.setProducerName(null); final NonPersistentPublisherStatsImpl publisherStats4 = new NonPersistentPublisherStatsImpl(); - publisherStats4.setSupportsPartialProducer(true); publisherStats4.setProducerName(null); topicStats2.addPublisher(publisherStats3); topicStats2.addPublisher(publisherStats4); - assertEquals(topicStats2.getPublishers().size(), 2); - // when the producerName is null, fall back to false - assertFalse(topicStats2.getPublishers().get(0).isSupportsPartialProducer()); - assertFalse(topicStats2.getPublishers().get(1).isSupportsPartialProducer()); + assertEquals(topicStats2.getPublishers().size(), 1); final NonPersistentPartitionedTopicStatsImpl target = new NonPersistentPartitionedTopicStatsImpl(); target.add(topicStats1); target.add(topicStats2); - assertEquals(target.getPublishers().size(), 2); + assertEquals(target.getPublishers().size(), 1); } + + @Test + public void testNonPersistentTopicStatsByDifferentPublishers() { + NonPersistentTopicStatsImpl topicStats = new NonPersistentTopicStatsImpl(); + NonPersistentTopicStatsImpl s1 = new NonPersistentTopicStatsImpl(); + NonPersistentTopicStatsImpl s2 = new NonPersistentTopicStatsImpl(); + NonPersistentPublisherStatsImpl p1 = new NonPersistentPublisherStatsImpl(); + NonPersistentPublisherStatsImpl p2 = new NonPersistentPublisherStatsImpl(); + NonPersistentPublisherStatsImpl p3 = new NonPersistentPublisherStatsImpl(); + p1.setProducerName("p1"); + p1.setMsgRateIn(1); + p2.setProducerName("p2"); + p2.setMsgRateIn(2); + p3.setMsgRateIn(3); + s1.addPublisher(p1); + s1.addPublisher(p2); + s1.addPublisher(p3); + s2.addPublisher(p1); + s2.addPublisher(p2); + topicStats.add(s1); + topicStats.add(s2); + + assertEquals(topicStats.getPublishers().size(), 3); + assertEquals(topicStats.getPublishers().get(0).getMsgRateIn(), 2); + assertEquals(topicStats.getPublishers().get(1).getMsgRateIn(), 4); + assertEquals(topicStats.getPublishers().get(2).getMsgRateIn(), 3); + } + + @Test + public void jsonWriteAndReadTest() throws JsonProcessingException { + ObjectMapper mapper = ObjectMapperFactory.create(); + + final NonPersistentPublisherStatsImpl publisherStats1 = new NonPersistentPublisherStatsImpl(); + publisherStats1.setProducerName("p1"); + final NonPersistentTopicStatsImpl src = new NonPersistentTopicStatsImpl(); + src.msgDropRate = 1.0; + src.addPublisher(publisherStats1); + + String json = mapper.writeValueAsString(src); + NonPersistentTopicStatsImpl dst = + (NonPersistentTopicStatsImpl) mapper.readValue(json, NonPersistentTopicStats.class); + assertEquals(dst.getPublishers().get(0).getProducerName(), "p1"); + assertEquals(dst.msgDropRate, 1.0); + } + } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PartitionedTopicStatsTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PartitionedTopicStatsTest.java index 9b7393d226a96..c64bd924b7534 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PartitionedTopicStatsTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PartitionedTopicStatsTest.java @@ -20,10 +20,14 @@ import static org.testng.Assert.assertEquals; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.pulsar.common.policies.data.stats.PartitionedTopicStatsImpl; import org.apache.pulsar.common.policies.data.stats.PublisherStatsImpl; import org.apache.pulsar.common.policies.data.stats.ReplicatorStatsImpl; import org.apache.pulsar.common.policies.data.stats.SubscriptionStatsImpl; +import org.apache.pulsar.common.policies.data.stats.TopicStatsImpl; +import org.apache.pulsar.common.util.ObjectMapperFactory; import org.testng.annotations.Test; public class PartitionedTopicStatsTest { @@ -55,4 +59,21 @@ public void testPartitionedTopicStats() { assertEquals(partitionedTopicStats.metadata.partitions, 0); assertEquals(partitionedTopicStats.partitions.size(), 0); } + + @Test + public void jsonWriteAndReadTest() throws JsonProcessingException { + ObjectMapper mapper = ObjectMapperFactory.create(); + + final PublisherStatsImpl publisherStats1 = new PublisherStatsImpl(); + publisherStats1.setProducerName("p1"); + final TopicStatsImpl src = new TopicStatsImpl(); + src.averageMsgSize = 1.0; + src.addPublisher(publisherStats1); + + String json = mapper.writeValueAsString(src); + TopicStatsImpl dst = (TopicStatsImpl) mapper.readValue(json, TopicStats.class); + assertEquals(dst.getPublishers().get(0).getProducerName(), "p1"); + assertEquals(dst.averageMsgSize, 1.0); + } + } \ No newline at end of file diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PersistentTopicStatsTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PersistentTopicStatsTest.java index 03bde68022f86..42360bf223a70 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PersistentTopicStatsTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PersistentTopicStatsTest.java @@ -19,7 +19,6 @@ package org.apache.pulsar.common.policies.data; import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertFalse; import org.apache.pulsar.common.policies.data.stats.PublisherStatsImpl; import org.apache.pulsar.common.policies.data.stats.ReplicatorStatsImpl; @@ -87,7 +86,6 @@ public void testPersistentTopicStatsAggregationPartialProducerIsNotSupported() { topicStats1.averageMsgSize = 1; topicStats1.storageSize = 1; final PublisherStatsImpl publisherStats1 = new PublisherStatsImpl(); - publisherStats1.setSupportsPartialProducer(false); publisherStats1.setProducerName("name1"); topicStats1.addPublisher(publisherStats1); topicStats1.subscriptions.put("test_ns", new SubscriptionStatsImpl()); @@ -101,7 +99,6 @@ public void testPersistentTopicStatsAggregationPartialProducerIsNotSupported() { topicStats2.averageMsgSize = 5; topicStats2.storageSize = 6; final PublisherStatsImpl publisherStats2 = new PublisherStatsImpl(); - publisherStats2.setSupportsPartialProducer(false); publisherStats2.setProducerName("name1"); topicStats2.addPublisher(publisherStats2); topicStats2.subscriptions.put("test_ns", new SubscriptionStatsImpl()); @@ -132,7 +129,6 @@ public void testPersistentTopicStatsAggregationPartialProducerSupported() { topicStats1.averageMsgSize = 1; topicStats1.storageSize = 1; final PublisherStatsImpl publisherStats1 = new PublisherStatsImpl(); - publisherStats1.setSupportsPartialProducer(true); publisherStats1.setProducerName("name1"); topicStats1.addPublisher(publisherStats1); topicStats1.subscriptions.put("test_ns", new SubscriptionStatsImpl()); @@ -146,7 +142,6 @@ public void testPersistentTopicStatsAggregationPartialProducerSupported() { topicStats2.averageMsgSize = 5; topicStats2.storageSize = 6; final PublisherStatsImpl publisherStats2 = new PublisherStatsImpl(); - publisherStats2.setSupportsPartialProducer(true); publisherStats2.setProducerName("name1"); topicStats2.addPublisher(publisherStats2); topicStats2.subscriptions.put("test_ns", new SubscriptionStatsImpl()); @@ -178,7 +173,6 @@ public void testPersistentTopicStatsAggregationByProducerName() { topicStats1.averageMsgSize = 1; topicStats1.storageSize = 1; final PublisherStatsImpl publisherStats1 = new PublisherStatsImpl(); - publisherStats1.setSupportsPartialProducer(true); publisherStats1.msgRateIn = 1; publisherStats1.setProducerName("name1"); topicStats1.addPublisher(publisherStats1); @@ -193,7 +187,6 @@ public void testPersistentTopicStatsAggregationByProducerName() { topicStats2.averageMsgSize = 5; topicStats2.storageSize = 6; final PublisherStatsImpl publisherStats2 = new PublisherStatsImpl(); - publisherStats2.setSupportsPartialProducer(true); publisherStats2.msgRateIn = 1; publisherStats2.setProducerName("name1"); topicStats2.addPublisher(publisherStats2); @@ -208,7 +201,6 @@ public void testPersistentTopicStatsAggregationByProducerName() { topicStats3.averageMsgSize = 0; topicStats3.storageSize = 0; final PublisherStatsImpl publisherStats3 = new PublisherStatsImpl(); - publisherStats3.setSupportsPartialProducer(true); publisherStats3.msgRateIn = 1; publisherStats3.setProducerName("name2"); topicStats3.addPublisher(publisherStats3); @@ -240,37 +232,55 @@ public void testPersistentTopicStatsAggregationByProducerName() { public void testPersistentTopicStatsByNullProducerName() { final TopicStatsImpl topicStats1 = new TopicStatsImpl(); final PublisherStatsImpl publisherStats1 = new PublisherStatsImpl(); - publisherStats1.setSupportsPartialProducer(false); publisherStats1.setProducerName(null); final PublisherStatsImpl publisherStats2 = new PublisherStatsImpl(); - publisherStats2.setSupportsPartialProducer(false); publisherStats2.setProducerName(null); topicStats1.addPublisher(publisherStats1); topicStats1.addPublisher(publisherStats2); - assertEquals(topicStats1.getPublishers().size(), 2); - assertFalse(topicStats1.getPublishers().get(0).isSupportsPartialProducer()); - assertFalse(topicStats1.getPublishers().get(1).isSupportsPartialProducer()); + assertEquals(topicStats1.getPublishers().size(), 1); final TopicStatsImpl topicStats2 = new TopicStatsImpl(); final PublisherStatsImpl publisherStats3 = new PublisherStatsImpl(); - publisherStats3.setSupportsPartialProducer(true); publisherStats3.setProducerName(null); final PublisherStatsImpl publisherStats4 = new PublisherStatsImpl(); - publisherStats4.setSupportsPartialProducer(true); publisherStats4.setProducerName(null); topicStats2.addPublisher(publisherStats3); topicStats2.addPublisher(publisherStats4); - assertEquals(topicStats2.getPublishers().size(), 2); - // when the producerName is null, fall back to false - assertFalse(topicStats2.getPublishers().get(0).isSupportsPartialProducer()); - assertFalse(topicStats2.getPublishers().get(1).isSupportsPartialProducer()); - + assertEquals(topicStats2.getPublishers().size(), 1); final TopicStatsImpl target = new TopicStatsImpl(); target.add(topicStats1); target.add(topicStats2); - assertEquals(target.getPublishers().size(), 2); + assertEquals(target.getPublishers().size(), 1); + } + + @Test + public void testPersistentTopicStatsByDifferentPublishers() { + TopicStatsImpl topicStats = new TopicStatsImpl(); + TopicStatsImpl s1 = new TopicStatsImpl(); + TopicStatsImpl s2 = new TopicStatsImpl(); + PublisherStatsImpl p1 = new PublisherStatsImpl(); + PublisherStatsImpl p2 = new PublisherStatsImpl(); + PublisherStatsImpl p3 = new PublisherStatsImpl(); + p1.setProducerName("p1"); + p1.setMsgRateIn(1); + p2.setProducerName("p2"); + p2.setMsgRateIn(2); + p3.setMsgRateIn(3); + s1.addPublisher(p1); + s1.addPublisher(p2); + s1.addPublisher(p3); + s2.addPublisher(p1); + s2.addPublisher(p2); + topicStats.add(s1); + topicStats.add(s2); + + assertEquals(topicStats.getPublishers().size(), 3); + assertEquals(topicStats.getPublishers().get(0).getMsgRateIn(), 2); + assertEquals(topicStats.getPublishers().get(1).getMsgRateIn(), 4); + assertEquals(topicStats.getPublishers().get(2).getMsgRateIn(), 3); + } } diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PublisherStatsTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PublisherStatsTest.java index 8be02c58eb0bd..28c8bcdbd6713 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PublisherStatsTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/policies/data/PublisherStatsTest.java @@ -54,7 +54,7 @@ public void testPublisherStats() throws Exception { assertNull(stats.getAddress()); assertNull(stats.getClientVersion()); assertNull(stats.getConnectedSince()); - assertNull(stats.getProducerName()); + assertEquals(stats.getProducerName(), "null"); stats.setAddress("address"); assertEquals(stats.getAddress(), "address"); @@ -95,7 +95,7 @@ public void testPublisherStats() throws Exception { assertEquals(stats.getClientVersion(), "version2"); stats.setProducerName(null); - assertNull(stats.getProducerName()); + assertEquals(stats.getProducerName(), "null"); assertNull(stats.getAddress());