From b9dec04154a9a4dc2dd6845236b10acfcdd9c1aa Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Wed, 19 Feb 2025 00:31:54 -0800 Subject: [PATCH] [fix][broker] Rate-Limiter fails with a huge spike in a traffic, and publish/consume stuck for a longer time --- .../pulsar/broker/qos/AsyncTokenBucket.java | 9 ++++--- .../broker/qos/AsyncTokenBucketTest.java | 24 +++++++++++++++++++ 2 files changed, 30 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/qos/AsyncTokenBucket.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/qos/AsyncTokenBucket.java index 8c43fa0a816fa..6fca1594999b9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/qos/AsyncTokenBucket.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/qos/AsyncTokenBucket.java @@ -193,12 +193,15 @@ private long consumeTokensAndMaybeUpdateTokensBalance(long consumeTokens, boolea // calculate the token delta by subtracting the consumed tokens from the new tokens long tokenDelta = newTokens - currentPendingConsumedTokens; if (tokenDelta != 0 || consumeTokens != 0) { + // prevent tokens to become excessive -ve where it can't recover + long cT = tokens < 0 ? 0 : consumeTokens; // update the tokens and return the current token value - return TOKENS_UPDATER.updateAndGet(this, + long availableTokens = TOKENS_UPDATER.updateAndGet(this, // limit the tokens to the capacity of the bucket - currentTokens -> Math.min(currentTokens + tokenDelta, getCapacity()) + currentTokens -> Math.min(currentTokens + Math.max(0, tokenDelta), getCapacity()) // subtract the consumed tokens from the capped tokens - - consumeTokens); + - cT); + return availableTokens; } else { return tokens; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/qos/AsyncTokenBucketTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/qos/AsyncTokenBucketTest.java index 82793f2748d78..c236122050e06 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/qos/AsyncTokenBucketTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/qos/AsyncTokenBucketTest.java @@ -21,6 +21,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertTrue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import org.testng.annotations.BeforeMethod; @@ -48,6 +49,29 @@ private void incrementMillis(long millis) { manualClockSource.addAndGet(TimeUnit.MILLISECONDS.toNanos(millis)); } + @Test + void testAsyncTokenWithMultiCall() throws Exception { + int rate = 2000; + int resolutionTimeNano = 8; + asyncTokenBucket = AsyncTokenBucket.builder().rate(rate).ratePeriodNanos(TimeUnit.SECONDS.toNanos(1)).clock( + new DefaultMonotonicSnapshotClock(TimeUnit.MILLISECONDS.toNanos(resolutionTimeNano), System::nanoTime)) + .build(); + + for (int i = 0; i < (1000); i++) { + for (int j = 0; j < (1000); j++) { + long token = asyncTokenBucket.getTokens(); + if (token < 0) { + // sleep to allow add new more tokens + Thread.sleep(resolutionTimeNano * 5); + assertTrue(asyncTokenBucket.getTokens() > 0); + } + // calling consumeTokens iteratively to simulate calling this method multiple times from multiple + // threads + asyncTokenBucket.consumeTokens(100); + } + } + } + @Test void shouldAddTokensWithConfiguredRate() { asyncTokenBucket =