Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions conf/broker.conf
Original file line number Diff line number Diff line change
Expand Up @@ -1279,6 +1279,10 @@ 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

### --- Schema storage --- ###
# The schema storage implementation used by this broker
schemaRegistryStorageClassName=org.apache.pulsar.broker.service.schema.BookkeeperSchemaStorageFactory
Expand Down
4 changes: 4 additions & 0 deletions conf/standalone.conf
Original file line number Diff line number Diff line change
Expand Up @@ -907,6 +907,10 @@ exposePreciseBacklogInPrometheus=false

splitTopicAndPartitionLabelInPrometheus=false

# If true, aggregate publisher stats of PartitionedTopicStats by producerName.
# Otherwise, aggregate it by list index.
aggregatePublisherStatsByProducerName=false

### --- Deprecated config variables --- ###

# Deprecated. Use configurationStoreServers
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2322,6 +2322,11 @@ public class ServiceConfiguration implements PulsarConfiguration {
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. --- ****/
/****
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLongFieldUpdater;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.service.BrokerServiceException.TopicClosedException;
import org.apache.pulsar.broker.service.BrokerServiceException.TopicTerminatedException;
import org.apache.pulsar.broker.service.Topic.PublishContext;
Expand Down Expand Up @@ -98,7 +99,10 @@ public Producer(Topic topic, TransportCnx cnx, long producerId, String producerN
boolean isEncrypted, Map<String, String> metadata, SchemaVersion schemaVersion, long epoch,
boolean userProvidedProducerName,
ProducerAccessMode accessMode,
Optional<Long> topicEpoch) {
Optional<Long> topicEpoch,
boolean supportsPartialProducer) {
final ServiceConfiguration serviceConf = cnx.getBrokerService().pulsar().getConfiguration();

this.topic = topic;
this.cnx = cnx;
this.producerId = producerId;
Expand All @@ -124,11 +128,20 @@ public Producer(Topic topic, TransportCnx cnx, long producerId, String producerN
stats.setClientVersion(cnx.getClientVersion());
stats.setProducerName(producerName);
stats.producerId = producerId;
if (serviceConf.isAggregatePublisherStatsByProducerName()) {
// 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);


String replicatorPrefix = cnx.getBrokerService().pulsar().getConfiguration().getReplicatorPrefix() + ".";
String replicatorPrefix = serviceConf.getReplicatorPrefix() + ".";
this.isRemote = producerName.startsWith(replicatorPrefix);
this.remoteCluster = parseRemoteClusterName(producerName, isRemote, replicatorPrefix);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1141,6 +1141,7 @@ protected void handleProducer(final CommandProducer cmdProducer) {
final Optional<Long> topicEpoch = cmdProducer.hasTopicEpoch()
? Optional.of(cmdProducer.getTopicEpoch()) : Optional.empty();
final boolean isTxnEnabled = cmdProducer.isTxnEnabled();
final boolean supportsPartialProducer = supportsPartialProducer();

TopicName topicName = validateTopicName(cmdProducer.getTopic(), requestId, cmdProducer);
if (topicName == null) {
Expand Down Expand Up @@ -1239,7 +1240,7 @@ protected void handleProducer(final CommandProducer cmdProducer) {
topic.checkIfTransactionBufferRecoverCompletely(isTxnEnabled).thenAccept(future -> {
buildProducerAndAddTopic(topic, producerId, producerName, requestId, isEncrypted,
metadata, schemaVersion, epoch, userProvidedProducerName, topicName,
producerAccessMode, topicEpoch, producerFuture);
producerAccessMode, topicEpoch, supportsPartialProducer, producerFuture);
}).exceptionally(exception -> {
Throwable cause = exception.getCause();
log.error("producerId {}, requestId {} : TransactionBuffer recover failed",
Expand Down Expand Up @@ -1305,11 +1306,12 @@ private void buildProducerAndAddTopic(Topic topic, long producerId, String produ
boolean isEncrypted, Map<String, String> metadata, SchemaVersion schemaVersion, long epoch,
boolean userProvidedProducerName, TopicName topicName,
ProducerAccessMode producerAccessMode,
Optional<Long> topicEpoch, CompletableFuture<Producer> producerFuture){
Optional<Long> topicEpoch, boolean supportsPartialProducer,
CompletableFuture<Producer> producerFuture){
CompletableFuture<Void> producerQueuedFuture = new CompletableFuture<>();
Producer producer = new Producer(topic, ServerCnx.this, producerId, producerName,
getPrincipal(), isEncrypted, metadata, schemaVersion, epoch,
userProvidedProducerName, producerAccessMode, topicEpoch);
userProvidedProducerName, producerAccessMode, topicEpoch, supportsPartialProducer);

topic.addProducer(producer, producerQueuedFuture).thenAccept(newTopicEpoch -> {
if (isActive()) {
Expand Down Expand Up @@ -2715,6 +2717,10 @@ boolean supportBrokerMetadata() {
return features != null && features.isSupportsBrokerEntryMetadata();
}

boolean supportsPartialProducer() {
return features != null && features.isSupportsPartialProducer();
}

@Override
public String getClientVersion() {
return clientVersion;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -824,7 +824,7 @@ public CompletableFuture<NonPersistentTopicStatsImpl> asyncGetStats(boolean getP
if (producer.isRemote()) {
remotePublishersStats.put(producer.getRemoteCluster(), publisherStats);
} else {
stats.getPublishers().add(publisherStats);
stats.addPublisher(publisherStats);
}
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1902,7 +1902,7 @@ public CompletableFuture<TopicStatsImpl> asyncGetStats(boolean getPreciseBacklog
if (producer.isRemote()) {
remotePublishersStats.put(producer.getRemoteCluster(), publisherStats);
} else {
stats.publishers.add(publisherStats);
stats.addPublisher(publisherStats);
}
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,14 +66,17 @@
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.MessageRouter;
import org.apache.pulsar.client.api.MessageRoutingMode;
import org.apache.pulsar.client.api.Producer;
import org.apache.pulsar.client.api.ProducerAccessMode;
import org.apache.pulsar.client.api.ProxyProtocol;
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.SubscriptionInitialPosition;
import org.apache.pulsar.client.api.SubscriptionType;
import org.apache.pulsar.client.api.TopicMetadata;
import org.apache.pulsar.client.impl.MessageIdImpl;
import org.apache.pulsar.common.naming.NamespaceName;
import org.apache.pulsar.common.naming.TopicDomain;
Expand Down Expand Up @@ -2222,7 +2225,6 @@ public void testGetTopicsWithDifferentMode() throws Exception {
producer2.close();
}


@Test(dataProvider = "isV1")
public void testNonPartitionedTopic(boolean isV1) throws Exception {
String tenant = "prop-xyz";
Expand Down Expand Up @@ -2270,4 +2272,83 @@ public void testFailedUpdatePartitionedTopic() throws Exception {
// validate subscription is created for new partition.
assertNotNull(admin.topics().getStats(partitionedTopicName + "-partition-" + 6).getSubscriptions().get(subName1));
}

@Test(dataProvider = "topicType")
public void testPartitionedStatsAggregationByProducerName(String topicType) throws Exception {
conf.setAggregatePublisherStatsByProducerName(true);
final String topic = topicType + "://prop-xyz/ns1/test-partitioned-stats-aggregation-by-producer-name";
admin.topics().createPartitionedTopic(topic, 10);

@Cleanup
Producer<byte[]> producer1 = pulsarClient.newProducer()
.topic(topic)
.enableLazyStartPartitionedProducers(true)
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.CustomPartition)
.messageRouter(new MessageRouter() {
@Override
public int choosePartition(Message<?> msg, TopicMetadata metadata) {
return msg.hasKey() ? Integer.parseInt(msg.getKey()) : 0;
}
})
.accessMode(ProducerAccessMode.Shared)
.create();

@Cleanup
Producer<byte[]> producer2 = pulsarClient.newProducer()
.topic(topic)
.enableLazyStartPartitionedProducers(true)
.enableBatching(false)
.messageRoutingMode(MessageRoutingMode.CustomPartition)
.messageRouter(new MessageRouter() {
@Override
public int choosePartition(Message<?> msg, TopicMetadata metadata) {
return msg.hasKey() ? Integer.parseInt(msg.getKey()) : 5;
}
})
.accessMode(ProducerAccessMode.Shared)
.create();

for (int i = 0; i < 10; i++) {
producer1.newMessage()
.key(String.valueOf(i % 5))
.value(("message".getBytes(StandardCharsets.UTF_8)))
.send();
producer2.newMessage()
.key(String.valueOf(i % 5 + 5))
.value(("message".getBytes(StandardCharsets.UTF_8)))
.send();
}

PartitionedTopicStats topicStats = admin.topics().getPartitionedStats(topic, true);
assertEquals(topicStats.getPartitions().size(), 10);
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);

@Cleanup
Producer<byte[]> producer1 = pulsarClient.newProducer()
.topic(topic + TopicName.PARTITIONED_TOPIC_SUFFIX + 0)
.create();

@Cleanup
Producer<byte[]> producer2 = pulsarClient.newProducer()
.topic(topic + TopicName.PARTITIONED_TOPIC_SUFFIX + 1)
.create();

PartitionedTopicStats topicStats = admin.topics().getPartitionedStats(topic, true);
assertEquals(topicStats.getPartitions().size(), 2);
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()));
}
}
Loading