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);