From 925264549e728e08f1ec09124254dfa6f4032426 Mon Sep 17 00:00:00 2001 From: penghui Date: Sun, 12 Dec 2021 13:08:27 +0800 Subject: [PATCH] Make sure the client can get response when encounter producer busy exception. When the producer with the same producer ID and same connection, Of course this usually doesn't happen. When cherry-picking #12846, the test using the same producer ID https://github.com/apache/pulsar/pull/12846/files#diff-b73845cd7706e03b89122a21b3e60c36c95ce6a0c7dd4574f8eda3fc67dd02b1R867-R876 And after we apply the change https://github.com/apache/pulsar/pull/8685/files#diff-1e0e8195fb5ec5a6d79acbc7d859c025a9b711f94e6ab37c94439e99b3202e84R1161 Only one request can get a response when encounter exceptions. The fix is make sure the broker don't miss the response to the client, we should make sure one request can get one response, otherwise the client side might get blocked. --- .../java/org/apache/pulsar/broker/service/ServerCnx.java | 7 +++---- .../org/apache/pulsar/broker/service/ServerCnxTest.java | 5 +++++ 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 65f756533aaf3..ac7c7b014bdc4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1322,10 +1322,9 @@ private void buildProducerAndAddTopic(Topic topic, long producerId, String produ remoteAddress, topicName, producerId, ex.getMessage()); producer.closeNow(true); - if (producerFuture.completeExceptionally(ex)) { - commandSender.sendErrorResponse(requestId, - BrokerServiceException.getClientErrorCode(ex), ex.getMessage()); - } + producerFuture.completeExceptionally(ex); + commandSender.sendErrorResponse(requestId, + BrokerServiceException.getClientErrorCode(ex), ex.getMessage()); return null; }); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java index 4fe14d89f7c67..67bbe0ae06d08 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java @@ -893,6 +893,11 @@ public void testCreateProducerTimeoutThenCreateSameNamedProducerShouldFail() thr producerName, Collections.emptyMap(), false); channel.writeInbound(createProducer3); + // 1nd producer will fail + response = getResponse(); + assertEquals(response.getClass(), CommandError.class); + assertEquals(((CommandError) response).getRequestId(), 1); + // 3nd producer will fail response = getResponse(); assertEquals(response.getClass(), CommandError.class);