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
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData;
import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException;
import org.apache.pulsar.broker.service.BrokerServiceException.ServerMetadataException;
import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter;
import org.apache.pulsar.client.impl.Murmur3Hash32;
import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType;
import org.apache.pulsar.common.util.FutureUtil;
Expand Down Expand Up @@ -261,6 +262,7 @@ public synchronized boolean canUnsubscribe(Consumer consumer) {
public CompletableFuture<Void> close(boolean disconnectConsumers,
Optional<BrokerLookupData> assignedBrokerLookupData) {
IS_CLOSED_UPDATER.set(this, TRUE);
getRateLimiter().ifPresent(DispatchRateLimiter::close);
Comment thread
lhotari marked this conversation as resolved.
return disconnectConsumers
? disconnectAllConsumers(false, assignedBrokerLookupData) : CompletableFuture.completedFuture(null);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.pulsar.broker.service;

import java.util.List;
Expand Down Expand Up @@ -101,7 +102,8 @@ default Optional<DispatchRateLimiter> getRateLimiter() {
}

default void updateRateLimiter() {
//No-op
initializeDispatchRateLimiterIfNeeded();
getRateLimiter().ifPresent(DispatchRateLimiter::updateDispatchRate);
}

default boolean initializeDispatchRateLimiterIfNeeded() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.pulsar.broker.service.RedeliveryTrackerDisabled;
import org.apache.pulsar.broker.service.SendMessageInfo;
import org.apache.pulsar.broker.service.Subscription;
import org.apache.pulsar.broker.service.persistent.DispatchRateLimiter;
import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.common.stats.Rate;
Expand Down Expand Up @@ -131,6 +132,7 @@ public synchronized boolean canUnsubscribe(Consumer consumer) {
public CompletableFuture<Void> close(boolean disconnectConsumers,
Optional<BrokerLookupData> assignedBrokerLookupData) {
IS_CLOSED_UPDATER.set(this, TRUE);
getRateLimiter().ifPresent(DispatchRateLimiter::close);
return disconnectConsumers
? disconnectAllConsumers(false, assignedBrokerLookupData) : CompletableFuture.completedFuture(null);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1026,13 +1026,6 @@ public Optional<DispatchRateLimiter> getRateLimiter() {
return dispatchRateLimiter;
}

@Override
public void updateRateLimiter() {
if (!initializeDispatchRateLimiterIfNeeded()) {
this.dispatchRateLimiter.ifPresent(DispatchRateLimiter::updateDispatchRate);
}
}

@Override
public boolean initializeDispatchRateLimiterIfNeeded() {
if (!dispatchRateLimiter.isPresent() && DispatchRateLimiter.isDispatchRateEnabled(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
Expand Down Expand Up @@ -560,13 +559,6 @@ public Optional<DispatchRateLimiter> getRateLimiter() {
return dispatchRateLimiter;
}

@Override
public void updateRateLimiter() {
if (!initializeDispatchRateLimiterIfNeeded()) {
this.dispatchRateLimiter.ifPresent(DispatchRateLimiter::updateDispatchRate);
}
}

@Override
public boolean initializeDispatchRateLimiterIfNeeded() {
if (!dispatchRateLimiter.isPresent() && DispatchRateLimiter.isDispatchRateEnabled(
Expand All @@ -578,13 +570,6 @@ public boolean initializeDispatchRateLimiterIfNeeded() {
return false;
}

@Override
public CompletableFuture<Void> close() {
IS_CLOSED_UPDATER.set(this, TRUE);
dispatchRateLimiter.ifPresent(DispatchRateLimiter::close);
return disconnectAllConsumers();
}

@Override
public boolean checkAndUnblockIfStuck() {
Consumer consumer = ACTIVE_CONSUMER_UPDATER.get(this);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@

import static org.awaitility.Awaitility.await;
import com.google.common.collect.Sets;
import java.time.Duration;
import java.util.Optional;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicInteger;
Expand Down Expand Up @@ -48,7 +47,7 @@ public class SubscriptionMessageDispatchThrottlingTest extends MessageDispatchTh
* @param subscription
* @throws Exception
*/
@Test(dataProvider = "subscriptionAndDispatchRateType", timeOut = 5000)
@Test(dataProvider = "subscriptionAndDispatchRateType", timeOut = 30000)
public void testMessageRateLimitingNotReceiveAllMessages(SubscriptionType subscription,
DispatchRateType dispatchRateType) throws Exception {
log.info("-- Starting {} test --", methodName);
Expand Down Expand Up @@ -144,7 +143,7 @@ public void testMessageRateLimitingNotReceiveAllMessages(SubscriptionType subscr
* @param subscription
* @throws Exception
*/
@Test(dataProvider = "subscriptions", timeOut = 5000)
@Test(dataProvider = "subscriptions", timeOut = 30000)
public void testMessageRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionType subscription)
throws Exception {
log.info("-- Starting {} test --", methodName);
Expand Down Expand Up @@ -218,7 +217,7 @@ public void testMessageRateLimitingReceiveAllMessagesAfterThrottling(Subscriptio
log.info("-- Exiting {} test --", methodName);
}

@Test(dataProvider = "subscriptions", timeOut = 30000, invocationCount = 15)
@Test(dataProvider = "subscriptions", timeOut = 30000)
private void testMessageNotDuplicated(SubscriptionType subscription) throws Exception {
int brokerRate = 1000;
int topicRate = 5000;
Expand Down Expand Up @@ -273,7 +272,7 @@ private void testMessageNotDuplicated(SubscriptionType subscription) throws Exce
Assert.fail("Should only have PersistentDispatcher in this test");
}
final DispatchRateLimiter subDispatchRateLimiter = subRateLimiter;
Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> {
Awaitility.await().untilAsserted(() -> {
DispatchRateLimiter brokerDispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter();
Assert.assertTrue(brokerDispatchRateLimiter != null
&& brokerDispatchRateLimiter.getDispatchRateOnByte() > 0);
Expand Down Expand Up @@ -320,7 +319,7 @@ private void testMessageNotDuplicated(SubscriptionType subscription) throws Exce
* @param subscription
* @throws Exception
*/
@Test(dataProvider = "subscriptions", timeOut = 5000)
@Test(dataProvider = "subscriptions", timeOut = 30000)
public void testBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionType subscription) throws Exception {
log.info("-- Starting {} test --", methodName);

Expand Down Expand Up @@ -451,7 +450,7 @@ private void testDispatchRate(SubscriptionType subscription,
Assert.fail("Should only have PersistentDispatcher in this test");
}
final DispatchRateLimiter subDispatchRateLimiter = subRateLimiter;
Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> {
Awaitility.await().untilAsserted(() -> {
DispatchRateLimiter brokerDispatchRateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter();
Assert.assertTrue(brokerDispatchRateLimiter != null
&& brokerDispatchRateLimiter.getDispatchRateOnByte() > 0);
Expand Down Expand Up @@ -526,7 +525,7 @@ public void testMultiLevelDispatch(SubscriptionType subscription) throws Excepti
* @param subscription
* @throws Exception
*/
@Test(dataProvider = "subscriptions", timeOut = 8000)
@Test(dataProvider = "subscriptions", timeOut = 30000)
public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(SubscriptionType subscription) throws Exception {
log.info("-- Starting {} test --", methodName);

Expand All @@ -539,6 +538,11 @@ public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(Subscri
long initBytes = pulsar.getConfiguration().getDispatchThrottlingRatePerTopicInByte();
final int byteRate = 1000;
admin.brokers().updateDynamicConfiguration("dispatchThrottlingRateInByte", "" + byteRate);

Awaitility.await().untilAsserted(() -> {
Assert.assertEquals(pulsar.getConfiguration().getDispatchThrottlingRateInByte(), byteRate);
});

admin.namespaces().createNamespace(namespace1, Sets.newHashSet("test"));
admin.namespaces().createNamespace(namespace2, Sets.newHashSet("test"));

Expand Down Expand Up @@ -572,7 +576,7 @@ public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(Subscri
Producer<byte[]> producer1 = pulsarClient.newProducer().topic(topicName1).create();
Producer<byte[]> producer2 = pulsarClient.newProducer().topic(topicName2).create();

Awaitility.await().atMost(Duration.ofMillis(500)).untilAsserted(() -> {
Awaitility.await().untilAsserted(() -> {
DispatchRateLimiter rateLimiter = pulsar.getBrokerService().getBrokerDispatchRateLimiter();
Assert.assertTrue(rateLimiter != null
&& rateLimiter.getDispatchRateOnByte() > 0);
Expand Down Expand Up @@ -605,7 +609,7 @@ public void testBrokerBytesRateLimitingReceiveAllMessagesAfterThrottling(Subscri
*
* @throws Exception
*/
@Test(timeOut = 5000)
@Test(timeOut = 30000)
public void testRateLimitingMultipleConsumers() throws Exception {
log.info("-- Starting {} test --", methodName);

Expand Down Expand Up @@ -691,7 +695,7 @@ public void testRateLimitingMultipleConsumers() throws Exception {
}


@Test(dataProvider = "subscriptions", timeOut = 5000)
@Test(dataProvider = "subscriptions", timeOut = 30000)
public void testClusterRateLimitingConfiguration(SubscriptionType subscription) throws Exception {
log.info("-- Starting {} test --", methodName);

Expand Down Expand Up @@ -868,7 +872,7 @@ public void testClusterPolicyOverrideConfiguration() throws Exception {
log.info("-- Exiting {} test --", methodName);
}

@Test(dataProvider = "subscriptions", timeOut = 11000)
@Test(dataProvider = "subscriptions", timeOut = 30000)
public void testClosingRateLimiter(SubscriptionType subscription) throws Exception {
log.info("-- Starting {} test --", methodName);

Expand Down