Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -18,23 +18,27 @@
*/
package org.apache.pulsar.broker.service;

import static com.google.common.base.Preconditions.checkArgument;
import static org.apache.bookkeeper.mledger.impl.ManagedLedgerMBeanImpl.ENTRY_LATENCY_BUCKETS_USEC;
import static org.apache.pulsar.broker.cache.ConfigurationCacheService.POLICIES;

import com.google.common.base.MoreObjects;
import java.util.Map;
import java.util.Objects;

import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.LongAdder;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.stream.Stream;

import com.google.common.collect.Sets;
import org.apache.bookkeeper.mledger.util.StatsBuckets;
import org.apache.pulsar.broker.admin.AdminResource;
import org.apache.pulsar.broker.service.schema.SchemaRegistryService;
import org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException;
import org.apache.pulsar.broker.stats.prometheus.metrics.Summary;
import org.apache.pulsar.common.api.proto.PulsarApi;
import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.common.policies.data.Policies;
import org.apache.pulsar.common.policies.data.PublishRate;
Expand All @@ -52,7 +56,7 @@ public abstract class AbstractTopic implements Topic {
protected final String topic;

// Producers currently connected to this topic
protected final ConcurrentHashMap<String, Producer> producers;
protected final ConcurrentHashMap<String, Set<Producer>> producerGroups;

protected final BrokerService brokerService;

Expand Down Expand Up @@ -90,7 +94,7 @@ public abstract class AbstractTopic implements Topic {
public AbstractTopic(String topic, BrokerService brokerService) {
this.topic = topic;
this.brokerService = brokerService;
this.producers = new ConcurrentHashMap<>();
this.producerGroups = new ConcurrentHashMap<>();
this.isFenced = false;
this.replicatorPrefix = brokerService.pulsar().getConfiguration().getReplicatorPrefix();
this.lastActive = System.nanoTime();
Expand All @@ -107,6 +111,26 @@ public AbstractTopic(String topic, BrokerService brokerService) {
updatePublishDispatcher(policies);
}

/**
* For testing purposes. Performs linear scan of all producers in the producer group.
*/
Producer getProducer(String producerName) {
return getProducer("", producerName);
}

/**
* For testing purposes. Performs linear scan of all producers in the producer group.
*/
Producer getProducer(String groupName, String producerName) {
Set<Producer> producers = producerGroups.get(groupName);
for (Producer producer : producers) {
if (producer.getProducerName().equals(producerName)) {
return producer;
}
}
return null;
}

protected boolean isProducersExceeded() {
Policies policies;
try {
Expand All @@ -119,21 +143,38 @@ protected boolean isProducersExceeded() {
final int maxProducers = policies.max_producers_per_topic > 0 ?
policies.max_producers_per_topic :
brokerService.pulsar().getConfiguration().getMaxProducersPerTopic();
if (maxProducers > 0 && maxProducers <= producers.size()) {
return true;
}
return false;
return maxProducers > 0 && maxProducers <= getProducers().count();
}

protected boolean hasLocalProducers() {
AtomicBoolean foundLocal = new AtomicBoolean(false);
producers.values().forEach(producer -> {
if (!producer.isRemote()) {
foundLocal.set(true);
@Override
public void removeProducer(Producer producer) {
checkArgument(producer.getTopic() == this);

boolean[] removed = { false };
producerGroups.computeIfPresent(producer.getProducerName(), (s, producerSet) -> {
if (producerSet instanceof HashSet) { // a non-exclusive producer group
if (producerSet.remove(producer)) {
removed[0] = true;
if (producerSet.size() == 0)
return null;
}
return producerSet;
} else { // an exclusive producer "group"
if (producerSet.contains(producer)) {
removed[0] = true;
return null;
} else {
return producerSet;
}
}
});
if (removed[0]) {
handleProducerRemoved(producer);
}
}

return foundLocal.get();
protected boolean hasLocalProducers() {
return getProducers().anyMatch(producer -> !producer.isRemote());
}

@Override
Expand All @@ -142,8 +183,8 @@ public String toString() {
}

@Override
public Map<String, Producer> getProducers() {
return producers;
public Stream<Producer> getProducers() {
return producerGroups.values().stream().flatMap(Collection::stream);
}


Expand Down Expand Up @@ -293,17 +334,17 @@ public void resetBrokerPublishCountAndEnableReadIfRequired(boolean doneBrokerRes
* it sets cnx auto-readable if producer's cnx is disabled due to publish-throttling
*/
protected void enableProducerReadForPublishRateLimiting() {
if (producers != null) {
producers.values().forEach(producer -> {
if (producerGroups != null) {
getProducers().forEach(producer -> {
producer.getCnx().cancelPublishRateLimiting();
producer.getCnx().enableCnxAutoRead();
});
}
}

protected void enableProducerReadForPublishBufferLimiting() {
if (producers != null) {
producers.values().forEach(producer -> {
if (producerGroups != null) {
getProducers().forEach(producer -> {
producer.getCnx().cancelPublishBufferLimiting();
producer.getCnx().enableCnxAutoRead();
});
Expand All @@ -327,32 +368,61 @@ protected void internalAddProducer(Producer producer) throws BrokerServiceExcept
log.debug("[{}] {} Got request to create producer ", topic, producer.getProducerName());
}

Producer existProducer = producers.putIfAbsent(producer.getProducerName(), producer);
if (existProducer != null) {
tryOverwriteOldProducer(existProducer, producer);
}
}

private void tryOverwriteOldProducer(Producer oldProducer, Producer newProducer)
throws BrokerServiceException {
boolean canOverwrite = false;
if (oldProducer.equals(newProducer) && !isUserProvidedProducerName(oldProducer)
&& !isUserProvidedProducerName(newProducer) && newProducer.getEpoch() > oldProducer.getEpoch()) {
oldProducer.close(false);
canOverwrite = true;
}
if (canOverwrite) {
if(!producers.replace(newProducer.getProducerName(), oldProducer, newProducer)) {
// Met concurrent update, throw exception here so that client can try reconnect later.
throw new BrokerServiceException.NamingException("Producer with name '" + newProducer.getProducerName()
+ "' replace concurrency error");
// the following the variables are used to get state out the compute function
// (we want to get out of producers.compute as quickly as possible so as not to block other concurrent actions)
BrokerServiceException[] bse = {null};
Collection<Producer> oldProducers = new LinkedList<>();
boolean parallelGroupMode = producer.getGroupMode() == PulsarApi.CommandProducer.GroupMode.Parallel;
producerGroups.compute(producer.getProducerName(), (s, producerSet) -> {
if (producerSet == null) { // no producer under that name and topic connected yet
if (parallelGroupMode) {
Set<Producer> identityHashSet = Sets.newIdentityHashSet();
identityHashSet.add(producer);
return identityHashSet;
} else {
return Collections.singleton(producer); // in reality a "set" that can contain only one element is enough here
}
} else {
handleProducerRemoved(oldProducer);
Producer existingProducer = producerSet.iterator().next();
boolean existingProducerIsExclusive =
existingProducer.getGroupMode() == PulsarApi.CommandProducer.GroupMode.Exclusive;
if (existingProducerIsExclusive) { // an exclusive producer is already connected under that producerName
if (parallelGroupMode) {
bse[0] = new BrokerServiceException.NamingException(
"Exclusive Producer with name '" + producer.getProducerName() + "' is already connected to topic");
return producerSet;
} else {
if (!isUserProvidedProducerName(existingProducer) && !isUserProvidedProducerName(producer)
&& producer.getEpoch() > existingProducer.getEpoch()) {
oldProducers.add(existingProducer);
return Collections.singleton(producer);
} else {
bse[0] = new BrokerServiceException.NamingException(
"Producer with name '" + producer.getProducerName() + "' is already connected to topic");
return producerSet;
}
}
} else { // a non-exclusive producer is already connected under that producerName
if (parallelGroupMode) {
if (!producerSet.add(producer)) {
bse[0] = new BrokerServiceException.NamingException(
"Non-exclusive producer with name '" + producer.getProducerName() + "' and address '" +
producer.getCnx().clientAddress() + "' is already connected to topic");
}
return producerSet;
} else {
oldProducers.addAll(producerSet);
return Collections.singleton(producer);
}
}
}
} else {
throw new BrokerServiceException.NamingException(
"Producer with name '" + newProducer.getProducerName() + "' is already connected to topic");
});
for (Producer oldProducer : oldProducers) {
oldProducer.close(false);
handleProducerRemoved(oldProducer);
}
if (bse[0] != null)
throw bse[0];
}

private boolean isUserProvidedProducerName(Producer producer){
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -184,9 +184,8 @@ private void dropBacklog(PersistentTopic persistentTopic, BacklogQuota quota) {
*/
private void disconnectProducers(PersistentTopic persistentTopic) {
List<CompletableFuture<Void>> futures = Lists.newArrayList();
Map<String, Producer> producers = persistentTopic.getProducers();

producers.values().forEach(producer -> {
persistentTopic.getProducers().forEach(producer -> {
log.info("Producer [{}] has exceeded backlog quota on topic [{}]. Disconnecting producer",
producer.getProducerName(), persistentTopic.getName());
futures.add(producer.disconnect());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2050,7 +2050,7 @@ private void checkMessagePublishBuffer() {
private void foreachProducer(Consumer<Producer> consumer) {
topics.forEach((n, t) -> {
Optional<Topic> topic = extractTopic(t);
topic.ifPresent(value -> value.getProducers().values().forEach(consumer));
topic.ifPresent(value -> value.getProducers().forEach(consumer));
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
import org.apache.pulsar.broker.service.Topic.PublishContext;
import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.common.api.proto.PulsarApi;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata;
import org.apache.pulsar.common.api.proto.PulsarApi.ServerError;
Expand All @@ -63,6 +64,7 @@ public class Producer {
private final String producerName;
private final long epoch;
private final boolean userProvidedProducerName;
private final PulsarApi.CommandProducer.GroupMode groupMode;
private final long producerId;
private final String appId;
private Rate msgIn;
Expand All @@ -88,13 +90,17 @@ public class Producer {
private final SchemaVersion schemaVersion;

public Producer(Topic topic, ServerCnx cnx, long producerId, String producerName, String appId,
boolean isEncrypted, Map<String, String> metadata, SchemaVersion schemaVersion, long epoch,
boolean userProvidedProducerName) {
boolean isEncrypted, Map<String, String> metadata, SchemaVersion schemaVersion, long epoch,
boolean userProvidedProducerName, PulsarApi.CommandProducer.GroupMode groupMode)
throws BrokerServiceException {
this.topic = topic;
this.cnx = cnx;
this.producerId = producerId;
this.producerName = checkNotNull(producerName);
this.userProvidedProducerName = userProvidedProducerName;
this.groupMode = groupMode;
if (groupMode != PulsarApi.CommandProducer.GroupMode.Exclusive && !userProvidedProducerName)
throw new BrokerServiceException.NotAllowedException("producerName must be specified in non-exclusive group modes");
this.epoch = epoch;
this.closeFuture = new CompletableFuture<>();
this.appId = appId;
Expand Down Expand Up @@ -277,6 +283,10 @@ public ServerCnx getCnx() {
return this.cnx;
}

public PulsarApi.CommandProducer.GroupMode getGroupMode() {
return groupMode;
}

private static final class MessagePublishContext implements PublishContext, Runnable {
private Producer producer;
private long sequenceId;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -929,6 +929,7 @@ protected void handleProducer(final CommandProducer cmdProducer) {
final boolean isEncrypted = cmdProducer.getEncrypted();
final Map<String, String> metadata = CommandUtils.metadataFromCommand(cmdProducer);
final SchemaData schema = cmdProducer.hasSchema() ? getSchema(cmdProducer.getSchema()) : null;
final CommandProducer.GroupMode groupMode = cmdProducer.getGroupMode();

TopicName topicName = validateTopicName(cmdProducer.getTopic(), requestId, cmdProducer);
if (topicName == null) {
Expand Down Expand Up @@ -1046,10 +1047,11 @@ protected void handleProducer(final CommandProducer cmdProducer) {
});

schemaVersionFuture.thenAccept(schemaVersion -> {
Producer producer = new Producer(topic, ServerCnx.this, producerId, producerName, authRole,
isEncrypted, metadata, schemaVersion, epoch, userProvidedProducerName);

try {
Producer producer = new Producer(topic, ServerCnx.this, producerId, producerName, authRole,
isEncrypted, metadata, schemaVersion, epoch, userProvidedProducerName,
groupMode);

topic.addProducer(producer);

if (isActive()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;

import org.apache.bookkeeper.mledger.Position;
import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter;
Expand Down Expand Up @@ -116,7 +117,7 @@ CompletableFuture<Subscription> createSubscription(String subscriptionName, Init

CompletableFuture<Void> delete();

Map<String, Producer> getProducers();
Stream<Producer> getProducers();

String getName();

Expand Down
Loading