Skip to content

Commit

Permalink
[improve] [proxy] Not close the socket if lookup failed caused by too…
Browse files Browse the repository at this point in the history
… many requests (apache#21216)

Motivation: The Pulsar client will close the socket if it receives a `ServiceNotReady` error when doing a lookup. The Broker will respond to the client with a `TooManyRequests` error if there are too many lookup requests in progress, but the Pulsar Proxy responds to the client with a `ServiceNotReady` error in the same scenario.

Modifications: Make Pulsar Proxy respond to the client with a `TooManyRequests` error if there are too many lookup requests in progress.
(cherry picked from commit d6c3fa4)
(cherry picked from commit 9a7c4bb)
  • Loading branch information
poorbarcode authored and mukesh-ctds committed Apr 15, 2024
1 parent bdd5e75 commit 3429c27
Show file tree
Hide file tree
Showing 3 changed files with 37 additions and 5 deletions.
4 changes: 0 additions & 4 deletions distribution/shell/src/assemble/LICENSE.bin.txt
Original file line number Diff line number Diff line change
Expand Up @@ -331,10 +331,6 @@ The Apache Software License, Version 2.0
- listenablefuture-9999.0-empty-to-avoid-conflict-with-guava.jar
* J2ObjC Annotations -- j2objc-annotations-1.3.jar
* Netty Reactive Streams -- netty-reactive-streams-2.0.6.jar
* Swagger
- swagger-annotations-1.6.10.jar
- swagger-core-1.6.10.jar
- swagger-models-1.6.10.jar
* DataSketches
- memory-0.8.3.jar
- sketches-core-0.8.3.jar
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ public void handleLookup(CommandLookupTopic lookup) {
log.debug("Lookup Request ID {} from {} rejected - {}.", clientRequestId, clientAddress,
throttlingErrorMessage);
}
writeAndFlush(Commands.newLookupErrorResponse(ServerError.ServiceNotReady,
writeAndFlush(Commands.newLookupErrorResponse(ServerError.TooManyRequests,
throttlingErrorMessage, clientRequestId));
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,18 +20,26 @@

import static org.mockito.Mockito.doReturn;
import static org.testng.Assert.assertTrue;
import static org.testng.Assert.fail;

import java.util.Optional;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;

import lombok.Cleanup;

import org.apache.pulsar.broker.BrokerTestUtil;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.broker.authentication.AuthenticationService;
import org.apache.pulsar.client.api.Producer;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.impl.BinaryProtoLookupService;
import org.apache.pulsar.client.impl.ClientCnx;
import org.apache.pulsar.client.impl.LookupService;
import org.apache.pulsar.client.impl.PulsarClientImpl;
import org.apache.pulsar.common.configuration.PulsarConfigurationLoader;
import org.apache.pulsar.common.naming.TopicName;
import org.apache.pulsar.metadata.impl.ZKMetadataStore;
import org.mockito.Mockito;
import org.testng.Assert;
Expand Down Expand Up @@ -112,4 +120,32 @@ public void testLookup() throws Exception {

Assert.assertEquals(LookupProxyHandler.REJECTED_PARTITIONS_METADATA_REQUESTS.get(), 5.0d);
}

@Test
public void testLookupThrottling() throws Exception {
PulsarClientImpl client = (PulsarClientImpl) PulsarClient.builder()
.serviceUrl(proxyService.getServiceUrl()).build();
String tpName = BrokerTestUtil.newUniqueName("persistent://public/default/tp");
LookupService lookupService = client.getLookup();
assertTrue(lookupService instanceof BinaryProtoLookupService);
ClientCnx lookupConnection = client.getCnxPool().getConnection(lookupService.resolveHost()).join();

// Make no permits to lookup.
Semaphore lookupSemaphore = proxyService.getLookupRequestSemaphore();
int availablePermits = lookupSemaphore.availablePermits();
lookupSemaphore.acquire(availablePermits);

// Verify will receive too many request exception, and the socket will not be closed.
try {
lookupService.getBroker(TopicName.get(tpName)).get();
fail("Expected too many request error.");
} catch (Exception ex) {
assertTrue(ex.getMessage().contains("Too many"));
}
assertTrue(lookupConnection.ctx().channel().isActive());

// cleanup.
lookupSemaphore.release(availablePermits);
client.close();
}
}

0 comments on commit 3429c27

Please sign in to comment.