diff --git a/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java b/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java index 7638bfa774d581..d790d091fc2792 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/cloud/rpc/MetaServiceProxyTest.java @@ -22,6 +22,7 @@ import org.apache.doris.common.jmockit.Deencapsulation; import org.apache.doris.rpc.RpcException; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.SettableFuture; import org.junit.After; import org.junit.Assert; @@ -315,11 +316,18 @@ public void testGetInstanceRateLimitedBeforeRpc() throws RpcException { @Test public void testGetVisibleVersionAsyncRateLimitedBeforeRpc() throws RpcException { - enableRateLimit(1, "", 1, 0); + // Consume limitForPeriod = qpsPerCore * CPU_CORES * burstSeconds in one weighted request. + // A long burst window also prevents a refresh before the assertion if the scheduler stalls. + int qpsPerCore = 1; + int rateLimitBurstSeconds = 60; + int limitForPeriod = qpsPerCore * CPU_CORES * rateLimitBurstSeconds; + enableRateLimit(qpsPerCore, "", rateLimitBurstSeconds, 0); MetaServiceProxy proxy = new MetaServiceProxy(); MetaServiceClient client = mockNormalClient(); putClient(proxy, client); - consumeRateLimitPermits(proxy, "getPartitionVersion"); + Mockito.when(client.getVisibleVersionAsync(Mockito.any())) + .thenReturn(Futures.immediateFuture(okGetVersionResponse())); + proxy.getVisibleVersionAsync(buildBatchPartitionVersionRequest(limitForPeriod)); try { proxy.getVisibleVersionAsync(Cloud.GetVersionRequest.newBuilder().build()); @@ -328,7 +336,7 @@ public void testGetVisibleVersionAsyncRateLimitedBeforeRpc() throws RpcException Assert.assertTrue(e.getMessage().contains("meta service rpc rate limited")); } - Mockito.verify(client, Mockito.never()).getVisibleVersionAsync(Mockito.any()); + Mockito.verify(client, Mockito.times(1)).getVisibleVersionAsync(Mockito.any()); Mockito.verify(client, Mockito.never()).shutdown(Mockito.anyBoolean()); } @@ -338,8 +346,8 @@ public void testBatchGetVisibleVersionAsyncConsumesMultipleRateLimitPermits() th MetaServiceProxy proxy = new MetaServiceProxy(); MetaServiceClient client = mockNormalClient(); putClient(proxy, client); - SettableFuture future = SettableFuture.create(); - Mockito.when(client.getVisibleVersionAsync(Mockito.any())).thenReturn(future); + Mockito.when(client.getVisibleVersionAsync(Mockito.any())) + .thenReturn(Futures.immediateFuture(okGetVersionResponse())); proxy.getVisibleVersionAsync(buildBatchTableVersionRequest(CPU_CORES));