From e45a34872baf839ad0522f83d3c555ad4c98d326 Mon Sep 17 00:00:00 2001 From: Michael Marshall Date: Tue, 25 Jan 2022 13:59:00 -0600 Subject: [PATCH] [ServerCnx] Only reply to client when completing producerFuture --- .../pulsar/broker/service/ServerCnx.java | 19 ++++++++++--------- 1 file changed, 10 insertions(+), 9 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 a605f312c6e22..7b5bb2b42ef68 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 @@ -1259,16 +1259,17 @@ protected void handleProducer(final CommandProducer cmdProducer) { (BrokerServiceException.TopicBacklogQuotaExceededException) cause; IllegalStateException illegalStateException = new IllegalStateException(tbqe); BacklogQuota.RetentionPolicy retentionPolicy = tbqe.getRetentionPolicy(); - if (retentionPolicy == BacklogQuota.RetentionPolicy.producer_request_hold) { - commandSender.sendErrorResponse(requestId, - ServerError.ProducerBlockedQuotaExceededError, - illegalStateException.getMessage()); - } else if (retentionPolicy == BacklogQuota.RetentionPolicy.producer_exception) { - commandSender.sendErrorResponse(requestId, - ServerError.ProducerBlockedQuotaExceededException, - illegalStateException.getMessage()); + if (producerFuture.completeExceptionally(illegalStateException)) { + if (retentionPolicy == BacklogQuota.RetentionPolicy.producer_request_hold) { + commandSender.sendErrorResponse(requestId, + ServerError.ProducerBlockedQuotaExceededError, + illegalStateException.getMessage()); + } else if (retentionPolicy == BacklogQuota.RetentionPolicy.producer_exception) { + commandSender.sendErrorResponse(requestId, + ServerError.ProducerBlockedQuotaExceededException, + illegalStateException.getMessage()); + } } - producerFuture.completeExceptionally(illegalStateException); producers.remove(producerId, producerFuture); return null; }