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 f4a8b750f84de..d8adcd0c42ec2 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 @@ -626,6 +626,62 @@ public class ServiceConfiguration implements PulsarConfiguration { + " non-backlog consumers as well.") private boolean dispatchThrottlingOnNonBacklogConsumerEnabled = false; + // <-- Pluggable dispatcher --> + @FieldContext( + category = CATEGORY_SERVER, + doc = "Path to a nar file containing all customized dispatchers if we want to use customized dispatcher." + ) + private String dispatcherNarPath = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for persistent Exclusive subscription mode." + ) + private String persistentDispatcherExclusive = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for persistent Shared subscription mode." + ) + private String persistentDispatcherShared = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for persistent Failover subscription mode." + ) + private String persistentDispatcherFailover = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for persistent Key_Shared subscription mode." + ) + private String persistentDispatcherKeyShared = null; + + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for non-persistent Exclusive subscription mode." + ) + private String nonpersistentDispatcherExclusive = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for non-persistent Shared subscription mode." + ) + private String nonpersistentDispatcherShared = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for non-persistent Failover subscription mode." + ) + private String nonpersistentDispatcherFailover = null; + + @FieldContext( + category = CATEGORY_SERVER, + doc = "Customized dispatcher to use for non-persistent Key_Shared subscription mode." + ) + private String nonpersistentDispatcherKeyShared = null; + // <-- dispatcher read settings --> @FieldContext( dynamic = true, diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index b43ec4d611a4a..13f9dd5517119 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -104,6 +104,8 @@ import org.apache.pulsar.broker.service.BrokerServiceException.PersistenceException; import org.apache.pulsar.broker.service.BrokerServiceException.ServerMetadataException; import org.apache.pulsar.broker.service.BrokerServiceException.ServiceUnitNotReadyException; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.DispatcherUtils; import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; @@ -446,6 +448,8 @@ public void start() throws Exception { // register listener to capture zk-latency ClientCnxnAspect.addListener(zkStatsListener); ClientCnxnAspect.registerExecutor(pulsar.getExecutor()); + + DispatcherUtils.init(serviceConfig.getDispatcherNarPath()); } protected void startStatsUpdater(int statsUpdateInitailDelayInSecs, int statsUpdateFrequencyInSecs) { @@ -680,6 +684,7 @@ public void close() throws IOException { ClientCnxnAspect.registerExecutor(null); topicOrderedExecutor.shutdown(); delayedDeliveryTrackerFactory.close(); + DispatcherUtils.close(); if (topicPublishRateLimiterMonitor != null) { topicPublishRateLimiterMonitor.shutdown(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 2762748c438b0..972431aea051b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -189,7 +189,7 @@ public String consumerName() { return consumerName; } - void notifyActiveConsumerChange(Consumer activeConsumer) { + public void notifyActiveConsumerChange(Consumer activeConsumer) { if (log.isDebugEnabled()) { log.debug("notify consumer {} - that [{}] for subscription {} has new active consumer : {}", consumerId, topicName, subscription.getName(), activeConsumer); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java index 8a0bd06641102..64dc3ace2d7f7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java @@ -24,6 +24,7 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.SubType; import org.apache.pulsar.common.api.proto.PulsarMarkers.ReplicatedSubscriptionsSnapshot; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractBaseDispatcher.java similarity index 96% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractBaseDispatcher.java index d0e2880d56a7f..4cd961d51deea 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractBaseDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractBaseDispatcher.java @@ -17,7 +17,7 @@ * under the License. */ -package org.apache.pulsar.broker.service; +package org.apache.pulsar.broker.service.dispatcher; import io.netty.buffer.ByteBuf; import java.io.IOException; @@ -30,6 +30,11 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.lang3.tuple.Pair; +import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; +import org.apache.pulsar.broker.service.EntryBatchSizes; +import org.apache.pulsar.broker.service.SendMessageInfo; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractDispatcherMultipleConsumers.java similarity index 98% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractDispatcherMultipleConsumers.java index 23b835efc2e07..f7b31f3bf406e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractDispatcherMultipleConsumers.java @@ -16,13 +16,16 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.service; +package org.apache.pulsar.broker.service.dispatcher; import com.carrotsearch.hppc.ObjectHashSet; import com.carrotsearch.hppc.ObjectSet; import java.util.Random; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; + +import org.apache.pulsar.broker.service.Consumer; +import org.apache.pulsar.broker.service.Subscription; import org.apache.pulsar.broker.service.persistent.PersistentStickyKeyDispatcherMultipleConsumers; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.SubType; import org.slf4j.Logger; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractDispatcherSingleActiveConsumer.java similarity index 96% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractDispatcherSingleActiveConsumer.java index 48eae73d345b4..8657b22ab2646 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/AbstractDispatcherSingleActiveConsumer.java @@ -16,7 +16,7 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.service; +package org.apache.pulsar.broker.service.dispatcher; import static com.google.common.base.Preconditions.checkArgument; import java.util.List; @@ -26,8 +26,14 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; + +import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; import org.apache.pulsar.broker.service.BrokerServiceException.ServerMetadataException; +import org.apache.pulsar.broker.service.Consumer; +import org.apache.pulsar.broker.service.HashRangeExclusiveStickyKeyConsumerSelector; +import org.apache.pulsar.broker.service.StickyKeyConsumerSelector; +import org.apache.pulsar.broker.service.Subscription; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.SubType; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/Dispatcher.java similarity index 91% rename from pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java rename to pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/Dispatcher.java index 9db773aa18ff7..068475ca49c8c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Dispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/Dispatcher.java @@ -16,12 +16,15 @@ * specific language governing permissions and limitations * under the License. */ -package org.apache.pulsar.broker.service; +package org.apache.pulsar.broker.service.dispatcher; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.pulsar.broker.service.BrokerServiceException; +import org.apache.pulsar.broker.service.Consumer; +import org.apache.pulsar.broker.service.RedeliveryTracker; import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.SubType; import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; @@ -119,4 +122,8 @@ default void markDeletePositionMoveForward() { // No-op } + default void init(DispatcherConfiguration dispatcherConfiguration) { + // No-op + } + } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherConfiguration.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherConfiguration.java new file mode 100644 index 0000000000000..7913dfd7cf195 --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherConfiguration.java @@ -0,0 +1,45 @@ +/** + * 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.broker.service.dispatcher; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.broker.service.Topic; +import org.apache.pulsar.common.api.proto.PulsarApi; + +/** + * Holds configuration for creating various kinds of dispatchers. + */ +@Data +@AllArgsConstructor +public class DispatcherConfiguration { + private PulsarApi.CommandSubscribe.SubType subType; + + private ManagedCursor cursor; + + private int partitionIndex; + + private Topic topic; + + private Subscription subscription; + + private PulsarApi.KeySharedMeta ksm; +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherFactory.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherFactory.java new file mode 100644 index 0000000000000..ceef7489bc65f --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherFactory.java @@ -0,0 +1,159 @@ +/** + * 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.broker.service.dispatcher; + +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.service.BrokerServiceException; +import org.apache.pulsar.broker.service.ConsistentHashingStickyKeyConsumerSelector; +import org.apache.pulsar.broker.service.HashRangeAutoSplitStickyKeyConsumerSelector; +import org.apache.pulsar.broker.service.HashRangeExclusiveStickyKeyConsumerSelector; +import org.apache.pulsar.broker.service.StickyKeyConsumerSelector; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentDispatcherSingleActiveConsumer; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentStickyKeyDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; +import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.persistent.PersistentDispatcherSingleActiveConsumer; +import org.apache.pulsar.broker.service.persistent.PersistentStickyKeyDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.common.api.proto.PulsarApi; + +import java.io.IOException; + +import static com.google.common.base.Preconditions.checkArgument; + +/** + * Factory class for creating {@link Dispatcher}. + */ +public class DispatcherFactory { + + /** + * Get proper dispatcher based on passed in config, could be default or customized dispatcher. + * @param dispatcherConfiguration Holds parameters for creating default dispatchers. + * @param serviceConfiguration For creating customized dispatcher. + * @return dispatcher + * @throws BrokerServiceException + */ + public static Dispatcher getDispatcher(DispatcherConfiguration dispatcherConfiguration, + ServiceConfiguration serviceConfiguration) + throws BrokerServiceException { + if (dispatcherConfiguration.getSubscription() instanceof PersistentSubscription) { + checkArgument(dispatcherConfiguration.getCursor() != null); + switch (dispatcherConfiguration.getSubType()) { + case Exclusive: + if (serviceConfiguration.getPersistentDispatcherExclusive() == null) { + return new PersistentDispatcherSingleActiveConsumer(dispatcherConfiguration.getCursor(), + PulsarApi.CommandSubscribe.SubType.Exclusive, 0, + (PersistentTopic) dispatcherConfiguration.getTopic(), dispatcherConfiguration.getSubscription()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getPersistentDispatcherExclusive(), dispatcherConfiguration); + } + case Shared: + if (serviceConfiguration.getPersistentDispatcherShared() == null) { + return new PersistentDispatcherMultipleConsumers((PersistentTopic) dispatcherConfiguration.getTopic(), + dispatcherConfiguration.getCursor(), dispatcherConfiguration.getSubscription()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getPersistentDispatcherShared(), dispatcherConfiguration); + } + case Failover: + if (serviceConfiguration.getPersistentDispatcherFailover() == null) { + return new PersistentDispatcherSingleActiveConsumer(dispatcherConfiguration.getCursor(), + PulsarApi.CommandSubscribe.SubType.Failover, dispatcherConfiguration.getPartitionIndex(), + (PersistentTopic) dispatcherConfiguration.getTopic(), dispatcherConfiguration.getSubscription()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getPersistentDispatcherFailover(), dispatcherConfiguration); + } + case Key_Shared: + checkArgument(dispatcherConfiguration.getKsm() != null); + if (serviceConfiguration.getPersistentDispatcherKeyShared() == null) { + return new PersistentStickyKeyDispatcherMultipleConsumers((PersistentTopic) dispatcherConfiguration.getTopic(), + dispatcherConfiguration.getCursor(), dispatcherConfiguration.getSubscription(), + serviceConfiguration, dispatcherConfiguration.getKsm()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getPersistentDispatcherKeyShared(), dispatcherConfiguration); + } + default: + throw new BrokerServiceException.ServerMetadataException("Unsupported subscription type"); + } + } else { + switch (dispatcherConfiguration.getSubType()) { + case Exclusive: + if (serviceConfiguration.getNonpersistentDispatcherExclusive() == null) { + return new NonPersistentDispatcherSingleActiveConsumer(PulsarApi.CommandSubscribe.SubType.Exclusive, + 0, (NonPersistentTopic) dispatcherConfiguration.getTopic(), + dispatcherConfiguration.getSubscription()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getNonpersistentDispatcherExclusive(), dispatcherConfiguration); + } + case Shared: + if (serviceConfiguration.getNonpersistentDispatcherShared() == null) { + return new NonPersistentDispatcherMultipleConsumers((NonPersistentTopic) dispatcherConfiguration.getTopic(), + dispatcherConfiguration.getSubscription()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getNonpersistentDispatcherShared(), dispatcherConfiguration); + } + case Failover: + if (serviceConfiguration.getNonpersistentDispatcherFailover() == null) { + return new NonPersistentDispatcherSingleActiveConsumer(PulsarApi.CommandSubscribe.SubType.Failover, + dispatcherConfiguration.getPartitionIndex(), + (NonPersistentTopic) dispatcherConfiguration.getTopic(), dispatcherConfiguration.getSubscription()); + } else { + return loadAndInitDispatcher(serviceConfiguration.getNonpersistentDispatcherFailover(), dispatcherConfiguration); + } + case Key_Shared: + checkArgument(dispatcherConfiguration.getKsm() != null); + if (serviceConfiguration.getNonpersistentDispatcherKeyShared() == null) { + switch (dispatcherConfiguration.getKsm().getKeySharedMode()) { + case STICKY: + return new NonPersistentStickyKeyDispatcherMultipleConsumers((NonPersistentTopic) dispatcherConfiguration.getTopic(), + dispatcherConfiguration.getSubscription(), new HashRangeExclusiveStickyKeyConsumerSelector()); + case AUTO_SPLIT: + default: + StickyKeyConsumerSelector selector; + if (serviceConfiguration.isSubscriptionKeySharedUseConsistentHashing()) { + selector = new ConsistentHashingStickyKeyConsumerSelector( + serviceConfiguration.getSubscriptionKeySharedConsistentHashingReplicaPoints()); + } else { + selector = new HashRangeAutoSplitStickyKeyConsumerSelector(); + } + + return new NonPersistentStickyKeyDispatcherMultipleConsumers((NonPersistentTopic) dispatcherConfiguration.getTopic(), + dispatcherConfiguration.getSubscription(), selector); + } + } else { + return loadAndInitDispatcher(serviceConfiguration.getNonpersistentDispatcherKeyShared(), dispatcherConfiguration); + } + default: + throw new BrokerServiceException.ServerMetadataException("Unsupported subscription type"); + } + } + } + + private static Dispatcher loadAndInitDispatcher(String dispatcherClassName, + DispatcherConfiguration dispatcherConfiguration) throws BrokerServiceException { + try { + Dispatcher dispatcher = DispatcherUtils.load(dispatcherClassName); + dispatcher.init(dispatcherConfiguration); + return dispatcher; + } catch (IOException e) { + throw new BrokerServiceException(e); + } + } +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherUtils.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherUtils.java new file mode 100644 index 0000000000000..62cf9cb6392af --- /dev/null +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/dispatcher/DispatcherUtils.java @@ -0,0 +1,117 @@ +/** + * 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.broker.service.dispatcher; + +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; +import org.apache.pulsar.common.nar.NarClassLoader; + +import java.io.File; +import java.io.IOException; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +import static com.google.common.base.Preconditions.checkArgument; + +/** + * Util class to load customized message dispatcher. + */ +@Slf4j +public class DispatcherUtils { + + // Classloader + private static NarClassLoader ncl = null; + + // Cache class definition. + private static Map dispatchers; + + /** + * Initialization. + * @param narPath Path for the nar containing customized dispatcher we need to load. + * @throws IOException + */ + public static void init(String narPath) throws IOException { + if (narPath != null) { + ncl = NarClassLoader.getFromArchive( + new File(narPath), + Collections.emptySet()); + dispatchers = new HashMap<>(); + } + } + + /** + * Load the customized {@link Dispatcher} according to the class name. + */ + public static Dispatcher load(String dispatcherClassName) throws IOException { + checkArgument(ncl != null); + if (StringUtils.isBlank(dispatcherClassName)) { + throw new IOException("Dispatcher class name can not be empty"); + } + + if (!dispatchers.containsKey(dispatcherClassName)) { + synchronized (DispatcherUtils.class) { + if (!dispatchers.containsKey(dispatcherClassName)) { + try { + Class dispatcherClazz = ncl.loadClass(dispatcherClassName); + dispatchers.put(dispatcherClassName, dispatcherClazz); + } catch (Throwable t) { + rethrowIOException(t); + } + } + } + } + Object dispatcher; + try { + dispatcher = dispatchers.get(dispatcherClassName).newInstance(); + if (!(dispatcher instanceof Dispatcher)) { + throw new RuntimeException("Class " + dispatcher.getClass() + + " does not implement dispatcher interface"); + } + return (Dispatcher) dispatcher; + } catch (Throwable t) { + rethrowIOException(t); + return null; + } + } + + public static void close() { + if (ncl != null) { + try { + ncl.close(); + } catch (IOException e) { + log.warn("Failed to close dispatcher class loader", e); + } + } + } + + private static void rethrowIOException(Throwable cause) + throws IOException { + if (cause instanceof IOException) { + throw (IOException) cause; + } else if (cause instanceof RuntimeException) { + throw (RuntimeException) cause; + } else if (cause instanceof Error) { + throw (Error) cause; + } else { + throw new IOException(cause.getMessage(), cause); + } + } + +} diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcher.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcher.java index 613a7b1a16598..3af7be330a7f3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcher.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcher.java @@ -24,7 +24,7 @@ import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.Consumer; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.SubType; import org.apache.pulsar.common.stats.Rate; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java index e222940c93a56..fad4b23314378 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherMultipleConsumers.java @@ -24,7 +24,7 @@ import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import org.apache.bookkeeper.mledger.Entry; import org.apache.pulsar.broker.ServiceConfiguration; -import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.dispatcher.AbstractDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; import org.apache.pulsar.broker.service.Consumer; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java index c77346df66bdb..90723913c0f86 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentDispatcherSingleActiveConsumer.java @@ -25,7 +25,7 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.admin.AdminResource; -import org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer; +import org.apache.pulsar.broker.service.dispatcher.AbstractDispatcherSingleActiveConsumer; import org.apache.pulsar.broker.service.Consumer; import org.apache.pulsar.broker.service.EntryBatchSizes; import org.apache.pulsar.broker.service.RedeliveryTracker; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentSubscription.java index a4d2d6f7ce9a5..7b21d2278adea 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentSubscription.java @@ -34,12 +34,14 @@ import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionFencedException; import org.apache.pulsar.broker.service.ConsistentHashingStickyKeyConsumerSelector; import org.apache.pulsar.broker.service.Consumer; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.HashRangeAutoSplitStickyKeyConsumerSelector; import org.apache.pulsar.broker.service.HashRangeExclusiveStickyKeyConsumerSelector; import org.apache.pulsar.broker.service.StickyKeyConsumerSelector; import org.apache.pulsar.broker.service.Subscription; import org.apache.pulsar.broker.service.Topic; +import org.apache.pulsar.broker.service.dispatcher.DispatcherConfiguration; +import org.apache.pulsar.broker.service.dispatcher.DispatcherFactory; import org.apache.pulsar.common.api.proto.PulsarApi.CommandAck.AckType; import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.SubType; import org.apache.pulsar.common.api.proto.PulsarApi.KeySharedMeta; @@ -100,18 +102,19 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce if (dispatcher == null || !dispatcher.isConsumerConnected()) { Dispatcher previousDispatcher = null; - + DispatcherConfiguration dispatcherConfiguration = new DispatcherConfiguration(consumer.subType(), + null, 0, topic, this, null); switch (consumer.subType()) { case Exclusive: if (dispatcher == null || dispatcher.getType() != SubType.Exclusive) { previousDispatcher = dispatcher; - dispatcher = new NonPersistentDispatcherSingleActiveConsumer(SubType.Exclusive, 0, topic, this); + dispatcher = (NonPersistentDispatcher) DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; case Shared: if (dispatcher == null || dispatcher.getType() != SubType.Shared) { previousDispatcher = dispatcher; - dispatcher = new NonPersistentDispatcherMultipleConsumers(topic, this); + dispatcher = (NonPersistentDispatcher) DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; case Failover: @@ -123,35 +126,16 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce if (dispatcher == null || dispatcher.getType() != SubType.Failover) { previousDispatcher = dispatcher; - dispatcher = new NonPersistentDispatcherSingleActiveConsumer(SubType.Failover, partitionIndex, - topic, this); + dispatcherConfiguration.setPartitionIndex(partitionIndex); + dispatcher = (NonPersistentDispatcher) DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; case Key_Shared: if (dispatcher == null || dispatcher.getType() != SubType.Key_Shared) { previousDispatcher = dispatcher; KeySharedMeta ksm = consumer.getKeySharedMeta() != null ? consumer.getKeySharedMeta() : KeySharedMeta.getDefaultInstance(); - - switch (ksm.getKeySharedMode()) { - case STICKY: - dispatcher = new NonPersistentStickyKeyDispatcherMultipleConsumers(topic, this, - new HashRangeExclusiveStickyKeyConsumerSelector()); - break; - - case AUTO_SPLIT: - default: - StickyKeyConsumerSelector selector; - ServiceConfiguration conf = topic.getBrokerService().getPulsar().getConfiguration(); - if (conf.isSubscriptionKeySharedUseConsistentHashing()) { - selector = new ConsistentHashingStickyKeyConsumerSelector( - conf.getSubscriptionKeySharedConsistentHashingReplicaPoints()); - } else { - selector = new HashRangeAutoSplitStickyKeyConsumerSelector(); - } - - dispatcher = new NonPersistentStickyKeyDispatcherMultipleConsumers(topic, this, selector); - break; - } + dispatcherConfiguration.setKsm(ksm); + dispatcher = (NonPersistentDispatcher) DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration());; } break; default: diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index f04c5126bc874..29dc3c1faff73 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -42,11 +42,11 @@ import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; -import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.dispatcher.AbstractDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; import org.apache.pulsar.broker.service.Consumer; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; import org.apache.pulsar.broker.service.EntryBatchSizes; import org.apache.pulsar.broker.service.InMemoryRedeliveryTracker; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 5d55b1dc5b3a1..11b8d7babb4c5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -38,9 +38,9 @@ import org.apache.bookkeeper.mledger.util.SafeRun; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.admin.AdminResource; -import org.apache.pulsar.broker.service.AbstractDispatcherSingleActiveConsumer; +import org.apache.pulsar.broker.service.dispatcher.AbstractDispatcherSingleActiveConsumer; import org.apache.pulsar.broker.service.Consumer; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.EntryBatchIndexesAcks; import org.apache.pulsar.broker.service.EntryBatchSizes; import org.apache.pulsar.broker.service.RedeliveryTracker; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index 1d5a031772d5d..8cb97f0362b0c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -67,7 +67,7 @@ public class PersistentStickyKeyDispatcherMultipleConsumers extends PersistentDi private final Set stuckConsumers; private final Set nextStuckConsumers; - PersistentStickyKeyDispatcherMultipleConsumers(PersistentTopic topic, ManagedCursor cursor, + public PersistentStickyKeyDispatcherMultipleConsumers(PersistentTopic topic, ManagedCursor cursor, Subscription subscription, ServiceConfiguration conf, KeySharedMeta ksm) { super(topic, cursor, subscription); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 491542a4edd48..52d229ff59212 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -51,9 +51,11 @@ import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionFencedException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionInvalidCursorPosition; import org.apache.pulsar.broker.service.Consumer; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.Subscription; import org.apache.pulsar.broker.service.Topic; +import org.apache.pulsar.broker.service.dispatcher.DispatcherConfiguration; +import org.apache.pulsar.broker.service.dispatcher.DispatcherFactory; import org.apache.pulsar.broker.transaction.pendingack.PendingAckHandle; import org.apache.pulsar.broker.transaction.pendingack.impl.PendingAckHandleDisabled; import org.apache.pulsar.broker.transaction.pendingack.impl.PendingAckHandleImpl; @@ -169,18 +171,20 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce if (dispatcher == null || !dispatcher.isConsumerConnected()) { Dispatcher previousDispatcher = null; + DispatcherConfiguration dispatcherConfiguration = new DispatcherConfiguration(consumer.subType(), + cursor, 0, topic, this, null); switch (consumer.subType()) { case Exclusive: if (dispatcher == null || dispatcher.getType() != SubType.Exclusive) { previousDispatcher = dispatcher; - dispatcher = new PersistentDispatcherSingleActiveConsumer(cursor, SubType.Exclusive, 0, topic, this); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; case Shared: if (dispatcher == null || dispatcher.getType() != SubType.Shared) { previousDispatcher = dispatcher; - dispatcher = new PersistentDispatcherMultipleConsumers(topic, cursor, this); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; case Failover: @@ -192,9 +196,9 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce } if (dispatcher == null || dispatcher.getType() != SubType.Failover) { + dispatcherConfiguration.setPartitionIndex(partitionIndex); previousDispatcher = dispatcher; - dispatcher = new PersistentDispatcherSingleActiveConsumer(cursor, SubType.Failover, partitionIndex, - topic, this); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; case Key_Shared: @@ -202,8 +206,8 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce previousDispatcher = dispatcher; KeySharedMeta ksm = consumer.getKeySharedMeta() != null ? consumer.getKeySharedMeta() : KeySharedMeta.getDefaultInstance(); - dispatcher = new PersistentStickyKeyDispatcherMultipleConsumers(topic, cursor, this, - topic.getBrokerService().getPulsar().getConfiguration(), ksm); + dispatcherConfiguration.setKsm(ksm); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, topic.getBrokerService().getPulsar().getConfiguration()); } break; default: diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index a88eb6f9baa26..1faee671196b6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -85,7 +85,7 @@ import org.apache.pulsar.broker.service.BrokerServiceException.TopicTerminatedException; import org.apache.pulsar.broker.service.BrokerServiceException.UnsupportedVersionException; import org.apache.pulsar.broker.service.Consumer; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.Producer; import org.apache.pulsar.broker.service.Replicator; import org.apache.pulsar.broker.service.StreamingStats; 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 6117c8f6e27d1..a353bb0c961f1 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 @@ -90,6 +90,7 @@ import org.apache.pulsar.broker.cache.ConfigurationCacheService; import org.apache.pulsar.broker.cache.LocalZooKeeperCacheService; import org.apache.pulsar.broker.namespace.NamespaceService; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.nonpersistent.NonPersistentReplicator; import org.apache.pulsar.broker.service.persistent.CompactorSubscription; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/DispatcherFactoryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/DispatcherFactoryTest.java new file mode 100644 index 0000000000000..889b516882e5d --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/DispatcherFactoryTest.java @@ -0,0 +1,116 @@ +/** + * 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.broker.service.dispatcher; + +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.service.BrokerService; +import org.apache.pulsar.broker.service.BrokerServiceException; +import org.apache.pulsar.broker.service.PersistentDispatcherFailoverConsumerTest; +import org.apache.pulsar.broker.service.Subscription; +import org.apache.pulsar.broker.service.Topic; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentDispatcherSingleActiveConsumer; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentStickyKeyDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentSubscription; +import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; +import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.persistent.PersistentDispatcherSingleActiveConsumer; +import org.apache.pulsar.broker.service.persistent.PersistentStickyKeyDispatcherMultipleConsumers; +import org.apache.pulsar.broker.service.persistent.PersistentSubscription; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; +import org.apache.pulsar.common.api.proto.PulsarApi; +import org.powermock.core.classloader.annotations.PowerMockIgnore; +import org.powermock.core.classloader.annotations.PrepareForTest; +import org.testng.annotations.Test; + +import static org.powermock.api.mockito.PowerMockito.mock; +import static org.powermock.api.mockito.PowerMockito.when; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; + +@PrepareForTest({ + PulsarApi.KeySharedMeta.class, DispatcherFactory.class +}) +@PowerMockIgnore({"org.apache.logging.log4j.*"}) +public class DispatcherFactoryTest { + + @Test + public void testCanGetDefaultDispatchers() throws BrokerServiceException { + ServiceConfiguration config = new ServiceConfiguration(); + ManagedCursor cursor = mock(ManagedCursor.class); + BrokerService brokerService = mock(BrokerService.class); + PulsarService pulsarService = mock(PulsarService.class); + PersistentTopic topic = mock(PersistentTopic.class); + when(pulsarService.getConfiguration()).thenReturn(config); + when(brokerService.pulsar()).thenReturn(pulsarService); + when(brokerService.getPulsar()).thenReturn(pulsarService); + when(topic.getBrokerService()).thenReturn(brokerService); + when(topic.getName()).thenReturn("my-topic"); + when(cursor.getName()).thenReturn("my-cursor"); + Subscription subscription = mock(PersistentSubscription.class); + PulsarApi.KeySharedMeta ksm = PulsarApi.KeySharedMeta.getDefaultInstance(); + + DispatcherConfiguration dispatcherConfiguration = new DispatcherConfiguration(PulsarApi.CommandSubscribe.SubType.Exclusive, + cursor, 0, topic, subscription, ksm); + Dispatcher dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof PersistentDispatcherSingleActiveConsumer); + assertEquals(dispatcher.getType(), PulsarApi.CommandSubscribe.SubType.Exclusive); + + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Failover); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + System.out.println(dispatcher.getClass()); + assertTrue(dispatcher instanceof PersistentDispatcherSingleActiveConsumer); + assertEquals(dispatcher.getType(), PulsarApi.CommandSubscribe.SubType.Failover); + + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Shared); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof PersistentDispatcherMultipleConsumers); + + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Key_Shared); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof PersistentStickyKeyDispatcherMultipleConsumers); + + NonPersistentTopic nonpersistentTopic = mock(NonPersistentTopic.class); + Subscription nonPersistentSubscription = mock(NonPersistentSubscription.class); + when(nonpersistentTopic.getBrokerService()).thenReturn(brokerService); + when(nonpersistentTopic.getName()).thenReturn("my-topic"); + dispatcherConfiguration.setTopic(nonpersistentTopic); + dispatcherConfiguration.setSubscription(nonPersistentSubscription); + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Exclusive); + + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof NonPersistentDispatcherSingleActiveConsumer); + assertEquals(dispatcher.getType(), PulsarApi.CommandSubscribe.SubType.Exclusive); + + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Failover); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof NonPersistentDispatcherSingleActiveConsumer); + assertEquals(dispatcher.getType(), PulsarApi.CommandSubscribe.SubType.Failover); + + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Shared); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof NonPersistentDispatcherMultipleConsumers); + + dispatcherConfiguration.setSubType(PulsarApi.CommandSubscribe.SubType.Key_Shared); + dispatcher = DispatcherFactory.getDispatcher(dispatcherConfiguration, config); + assertTrue(dispatcher instanceof NonPersistentStickyKeyDispatcherMultipleConsumers); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/DispatcherUtilsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/DispatcherUtilsTest.java new file mode 100644 index 0000000000000..d475c7d3b882a --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/DispatcherUtilsTest.java @@ -0,0 +1,96 @@ +/** + * 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.broker.service.dispatcher; + +import org.apache.pulsar.common.nar.NarClassLoader; +import org.powermock.api.mockito.PowerMockito; +import org.powermock.core.classloader.annotations.PowerMockIgnore; +import org.powermock.core.classloader.annotations.PrepareForTest; +import org.testng.IObjectFactory; +import org.testng.annotations.ObjectFactory; +import org.testng.annotations.Test; + +import java.io.File; +import java.util.Set; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; +import static org.testng.AssertJUnit.assertTrue; +import static org.testng.AssertJUnit.fail; + +@PrepareForTest({ + DispatcherUtils.class, NarClassLoader.class +}) +@PowerMockIgnore({"org.apache.logging.log4j.*"}) +public class DispatcherUtilsTest { + + // Necessary to make PowerMockito.mockStatic work with TestNG. + @ObjectFactory + public IObjectFactory getObjectFactory() { + return new org.powermock.modules.testng.PowerMockObjectFactory(); + } + + @Test + public void testLoadDispatcher() throws Exception { + String dispatcherClassName = "org.apache.pulsar.my.smart.dispatcher"; + String narPath = "/path/to/dispatcher/nar"; + + NarClassLoader mockLoader = mock(NarClassLoader.class); + Class dispatcherClass = MockCustomizedDispatcher.class; + when(mockLoader.loadClass(eq(dispatcherClassName))) + .thenReturn(dispatcherClass); + + PowerMockito.mockStatic(NarClassLoader.class); + PowerMockito.when(NarClassLoader.getFromArchive( + any(File.class), + any(Set.class) + )).thenReturn(mockLoader); + + DispatcherUtils.init(narPath); + Dispatcher returnedDispatcher = DispatcherUtils.load(dispatcherClassName); + + assertTrue(returnedDispatcher instanceof MockCustomizedDispatcher); + } + + @Test + public void testExceptionIfNarPathNotSpecified() throws Exception { + String dispatcherClassName = "org.apache.pulsar.my.smart.dispatcher"; + + NarClassLoader mockLoader = mock(NarClassLoader.class); + Class dispatcherClass = MockCustomizedDispatcher.class; + when(mockLoader.loadClass(eq(dispatcherClassName))) + .thenReturn(dispatcherClass); + + PowerMockito.mockStatic(NarClassLoader.class); + PowerMockito.when(NarClassLoader.getFromArchive( + any(File.class), + any(Set.class) + )).thenReturn(mockLoader); + + DispatcherUtils.init(null); + try { + DispatcherUtils.load(dispatcherClassName); + fail("Should not reach here"); + } catch (Exception e) { + assertTrue(e instanceof IllegalArgumentException); + } + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/MockCustomizedDispatcher.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/MockCustomizedDispatcher.java new file mode 100644 index 0000000000000..0a7547857cd81 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/dispatcher/MockCustomizedDispatcher.java @@ -0,0 +1,83 @@ +/** + * 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.broker.service.dispatcher; + +import lombok.NoArgsConstructor; +import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.pulsar.broker.service.BrokerServiceException; +import org.apache.pulsar.broker.service.Consumer; +import org.apache.pulsar.broker.service.RedeliveryTracker; +import org.apache.pulsar.common.api.proto.PulsarApi; + +import java.util.List; +import java.util.concurrent.CompletableFuture; + +@NoArgsConstructor +public class MockCustomizedDispatcher implements Dispatcher{ + @Override + public void addConsumer(Consumer consumer) throws BrokerServiceException { } + + @Override + public void removeConsumer(Consumer consumer) throws BrokerServiceException { } + + @Override + public void consumerFlow(Consumer consumer, int additionalNumberOfMessages) { } + + @Override + public boolean isConsumerConnected() { return false; } + + @Override + public List getConsumers() { return null; } + + @Override + public boolean canUnsubscribe(Consumer consumer) { return false; } + + @Override + public CompletableFuture close() { return null; } + + @Override + public boolean isClosed() { return false; } + + @Override + public CompletableFuture disconnectActiveConsumers(boolean isResetCursor) { return null; } + + @Override + public CompletableFuture disconnectAllConsumers(boolean isResetCursor) { return null; } + + @Override + public void resetCloseFuture() { } + + @Override + public void reset() { } + + @Override + public PulsarApi.CommandSubscribe.SubType getType() { return null; } + + @Override + public void redeliverUnacknowledgedMessages(Consumer consumer) { } + + @Override + public void redeliverUnacknowledgedMessages(Consumer consumer, List positions) { } + + @Override + public void addUnAckedMessages(int unAckMessages) { } + + @Override + public RedeliveryTracker getRedeliveryTracker() { return null; } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java index 799d7f1123612..1605d71e5eb3f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java @@ -34,7 +34,7 @@ import lombok.Cleanup; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java index 7ed88f5c516f2..11b79420b2fb8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SubscriptionMessageDispatchThrottlingTest.java @@ -23,7 +23,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicInteger; -import org.apache.pulsar.broker.service.Dispatcher; +import org.apache.pulsar.broker.service.dispatcher.Dispatcher; import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentDispatcherSingleActiveConsumer;