From 736333436ed5f95ec9696be303460549797f8735 Mon Sep 17 00:00:00 2001 From: Stanislav Kozlovski Date: Tue, 26 Feb 2019 11:16:06 +0000 Subject: [PATCH 1/3] Add teardown helper for static variables in GroupedUserQuotaCallback --- .../integration/kafka/api/CustomQuotaCallbackTest.scala | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala b/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala index fed5380afee4a..d2d2647a012f6 100644 --- a/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala +++ b/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala @@ -81,6 +81,7 @@ class CustomQuotaCallbackTest extends IntegrationTestHarness with SaslSetup { @After override def tearDown(): Unit = { adminClients.foreach(_.close()) + GroupedUserQuotaCallback.tearDown() super.tearDown() } @@ -321,6 +322,12 @@ object GroupedUserQuotaCallback { ClientQuotaType.REQUEST -> new AtomicInteger ) val callbackInstances = new AtomicInteger + + def tearDown(): Unit = { + callbackInstances.set(0) + quotaLimitCalls.values.foreach(_.set(0)) + UnlimitedQuotaMetricTags.clear() + } } /** From 23613f2700f9519c3fbb52cb30a3b3b624f82427 Mon Sep 17 00:00:00 2001 From: Stanislav Kozlovski Date: Tue, 26 Feb 2019 15:50:31 +0000 Subject: [PATCH 2/3] Add sleep after initial user quota configuration --- .../scala/integration/kafka/api/CustomQuotaCallbackTest.scala | 1 + 1 file changed, 1 insertion(+) diff --git a/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala b/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala index d2d2647a012f6..77f81c8200a8b 100644 --- a/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala +++ b/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala @@ -97,6 +97,7 @@ class CustomQuotaCallbackTest extends IntegrationTestHarness with SaslSetup { var brokerId = 0 var user = createGroupWithOneUser("group0_user1", brokerId) user.configureAndWaitForQuota(1000000, 2000000) + Thread.sleep(1000) // ensure brokers fully process dynamic config changes quotaLimitCalls.values.foreach(_.set(0)) user.produceConsume(expectProduceThrottle = false, expectConsumeThrottle = false) From 9e21360abcdfbfd2f1b405fcdace9dc0b7f54c49 Mon Sep 17 00:00:00 2001 From: Stanislav Kozlovski Date: Wed, 27 Feb 2019 09:39:20 +0000 Subject: [PATCH 3/3] Replace Thread.sleep with quick sanity check --- .../scala/integration/kafka/api/CustomQuotaCallbackTest.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala b/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala index 77f81c8200a8b..8c1d34ddc2c98 100644 --- a/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala +++ b/core/src/test/scala/integration/kafka/api/CustomQuotaCallbackTest.scala @@ -97,14 +97,13 @@ class CustomQuotaCallbackTest extends IntegrationTestHarness with SaslSetup { var brokerId = 0 var user = createGroupWithOneUser("group0_user1", brokerId) user.configureAndWaitForQuota(1000000, 2000000) - Thread.sleep(1000) // ensure brokers fully process dynamic config changes quotaLimitCalls.values.foreach(_.set(0)) user.produceConsume(expectProduceThrottle = false, expectConsumeThrottle = false) // ClientQuotaCallback#quotaLimit is invoked by each quota manager once for each new client assertEquals(1, quotaLimitCalls(ClientQuotaType.PRODUCE).get) assertEquals(1, quotaLimitCalls(ClientQuotaType.FETCH).get) - assertTrue(s"Too many quotaLimit calls $quotaLimitCalls", quotaLimitCalls(ClientQuotaType.REQUEST).get <= serverCount) + assertTrue(s"Too many quotaLimit calls $quotaLimitCalls", quotaLimitCalls(ClientQuotaType.REQUEST).get <= 10) // sanity check // Large quota updated to small quota, should throttle user.configureAndWaitForQuota(9000, 3000) user.produceConsume(expectProduceThrottle = true, expectConsumeThrottle = true)