From 9f8a5996c9a241b768b8b7cd3b9197552b764f97 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Thu, 21 Jan 2021 20:00:30 +0800 Subject: [PATCH 1/3] remove consumer unnecessary locks --- .../pulsar/client/impl/ConsumerBase.java | 22 ++-- .../pulsar/client/impl/ConsumerImpl.java | 105 +++++++----------- .../pulsar/client/impl/MessagesImpl.java | 12 ++ .../client/impl/MultiTopicsConsumerImpl.java | 68 +++++------- 4 files changed, 88 insertions(+), 119 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java index dfce6ed63335b..33b721c27683c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java @@ -35,8 +35,6 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLongFieldUpdater; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; import io.netty.util.Timeout; import org.apache.pulsar.client.api.BatchReceivePolicy; @@ -84,7 +82,6 @@ public abstract class ConsumerBase extends HandlerState implements Consumer conf, int receiverQueueSize, ExecutorProvider executorProvider, @@ -733,12 +730,8 @@ protected void notifyPendingBatchReceivedCallBack() { if (opBatchReceive == null) { return; } - try { - reentrantLock.lock(); - notifyPendingBatchReceivedCallBack(opBatchReceive); - } finally { - reentrantLock.unlock(); - } + + notifyPendingBatchReceivedCallBack(opBatchReceive); } private OpBatchReceive peekNextBatchReceive() { @@ -780,16 +773,17 @@ private OpBatchReceive pollNextBatchReceive() { protected final void notifyPendingBatchReceivedCallBack(OpBatchReceive opBatchReceive) { MessagesImpl messages = getNewMessagesImpl(); - Message msgPeeked = incomingMessages.peek(); - while (msgPeeked != null && messages.canAdd(msgPeeked)) { - Message msg = incomingMessages.poll(); + + Message msg = null; + do { + msg = incomingMessages.poll(); if (msg != null) { messageProcessed(msg); Message interceptMsg = beforeConsume(msg); messages.add(interceptMsg); } - msgPeeked = incomingMessages.peek(); - } + } while (msg != null && messages.canAdd()); + completePendingBatchReceive(opBatchReceive.future, messages); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index c7f67f48aef48..fb9bc624bee9c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -126,8 +126,6 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final int receiverQueueRefillThreshold; - private final ReadWriteLock lock = new ReentrantReadWriteLock(); - private final UnAckedMessageTracker unAckedMessageTracker; private final AcknowledgmentsGroupingTracker acknowledgmentsGroupingTracker; private final NegativeAcksTracker negativeAcksTracker; @@ -409,7 +407,6 @@ protected CompletableFuture> internalReceiveAsync() { CompletableFuture> result = cancellationHandler.createFuture(); Message message = null; try { - lock.writeLock().lock(); message = incomingMessages.poll(0, TimeUnit.MILLISECONDS); if (message == null) { pendingReceives.add(result); @@ -418,8 +415,6 @@ protected CompletableFuture> internalReceiveAsync() { } catch (InterruptedException e) { Thread.currentThread().interrupt(); result.completeExceptionally(e); - } finally { - lock.writeLock().unlock(); } if (message != null) { @@ -470,32 +465,28 @@ protected Messages internalBatchReceive() throws PulsarClientException { protected CompletableFuture> internalBatchReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - try { - lock.writeLock().lock(); - if (pendingBatchReceives == null) { - pendingBatchReceives = Queues.newConcurrentLinkedQueue(); - } - if (hasEnoughMessagesForBatchReceive()) { - MessagesImpl messages = getNewMessagesImpl(); - Message msgPeeked = incomingMessages.peek(); - while (msgPeeked != null && messages.canAdd(msgPeeked)) { - Message msg = incomingMessages.poll(); - if (msg != null) { - messageProcessed(msg); - Message interceptMsg = beforeConsume(msg); - messages.add(interceptMsg); - } - msgPeeked = incomingMessages.peek(); + if (pendingBatchReceives == null) { + pendingBatchReceives = Queues.newConcurrentLinkedQueue(); + } + if (hasEnoughMessagesForBatchReceive()) { + MessagesImpl messages = getNewMessagesImpl(); + Message msg = null; + do { + msg = incomingMessages.poll(); + if (msg != null) { + messageProcessed(msg); + Message interceptMsg = beforeConsume(msg); + messages.add(interceptMsg); } - result.complete(messages); - } else { - OpBatchReceive opBatchReceive = OpBatchReceive.of(result); - pendingBatchReceives.add(opBatchReceive); - cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); - } - } finally { - lock.writeLock().unlock(); + } while (msg != null && messages.canAdd()); + + result.complete(messages); + } else { + OpBatchReceive opBatchReceive = OpBatchReceive.of(result); + pendingBatchReceives.add(opBatchReceive); + cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); } + return result; } @@ -948,14 +939,9 @@ private void closeConsumerTasks() { } private void failPendingReceive() { - lock.readLock().lock(); - try { - if (pinnedExecutor != null && !pinnedExecutor.isShutdown()) { - failPendingReceives(this.pendingReceives); - failPendingBatchReceives(this.pendingBatchReceives); - } - } finally { - lock.readLock().unlock(); + if (pinnedExecutor != null && !pinnedExecutor.isShutdown()) { + failPendingReceives(this.pendingReceives); + failPendingBatchReceives(this.pendingBatchReceives); } } @@ -1053,23 +1039,18 @@ void messageReceived(MessageIdData messageId, int redeliveryCount, List ac uncompressedPayload, createEncryptionContext(msgMetadata), cnx, schema, redeliveryCount); uncompressedPayload.release(); - lock.readLock().lock(); - try { - // Enqueue the message so that it can be retrieved when application calls receive() - // if the conf.getReceiverQueueSize() is 0 then discard message if no one is waiting for it. - // if asyncReceive is waiting then notify callback without adding to incomingMessages queue - if (deadLetterPolicy != null && possibleSendToDeadLetterTopicMessages != null && redeliveryCount >= deadLetterPolicy.getMaxRedeliverCount()) { - possibleSendToDeadLetterTopicMessages.put((MessageIdImpl)message.getMessageId(), Collections.singletonList(message)); - } - if (peekPendingReceive() != null) { - notifyPendingReceivedCallback(message, null); - } else if (enqueueMessageAndCheckBatchReceive(message)) { - if (hasPendingBatchReceive()) { - notifyPendingBatchReceivedCallBack(); - } - } - } finally { - lock.readLock().unlock(); + // Enqueue the message so that it can be retrieved when application calls receive() + // if the conf.getReceiverQueueSize() is 0 then discard message if no one is waiting for it. + // if asyncReceive is waiting then notify callback without adding to incomingMessages queue + if (deadLetterPolicy != null && possibleSendToDeadLetterTopicMessages != null && + redeliveryCount >= deadLetterPolicy.getMaxRedeliverCount()) { + possibleSendToDeadLetterTopicMessages.put((MessageIdImpl)message.getMessageId(), + Collections.singletonList(message)); + } + if (peekPendingReceive() != null) { + notifyPendingReceivedCallback(message, null); + } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { + notifyPendingBatchReceivedCallBack(); } } else { // handle batch message enqueuing; uncompressed payload has all messages in batch @@ -1280,17 +1261,11 @@ void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata, int redeliv if (possibleToDeadLetter != null) { possibleToDeadLetter.add(message); } - lock.readLock().lock(); - try { - if (peekPendingReceive() != null) { - notifyPendingReceivedCallback(message, null); - } else if (enqueueMessageAndCheckBatchReceive(message)) { - if (hasPendingBatchReceive()) { - notifyPendingBatchReceivedCallBack(); - } - } - } finally { - lock.readLock().unlock(); + + if (peekPendingReceive() != null) { + notifyPendingReceivedCallback(message, null); + } else if (enqueueMessageAndCheckBatchReceive(message) && hasPendingBatchReceive()) { + notifyPendingBatchReceivedCallBack(); } singleMessagePayload.release(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java index c56694e2759ea..f2705958784ef 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java @@ -56,6 +56,18 @@ protected boolean canAdd(Message message) { return true; } + protected boolean canAdd() { + if (maxNumberOfMessages > 0 && currentNumberOfMessages + 1 > maxNumberOfMessages) { + return false; + } + + if (maxSizeOfMessages > 0 && currentSizeOfMessages + 1 > maxSizeOfMessages) { + return false; + } + + return true; + } + protected void add(Message message) { if (message == null) { return; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java index 39f77157634d6..4c820dd680e31 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java @@ -62,8 +62,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.locks.ReadWriteLock; -import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -98,7 +96,6 @@ public class MultiTopicsConsumerImpl extends ConsumerBase { TopicsPartitionChangedListener topicsPartitionChangedListener; CompletableFuture partitionsAutoUpdateFuture = null; - private final ReadWriteLock lock = new ReentrantReadWriteLock(); private final ConsumerStatsRecorder stats; private final UnAckedMessageTracker unAckedMessageTracker; private final ConsumerConfigurationData internalConfig; @@ -359,33 +356,28 @@ protected Messages internalBatchReceive() throws PulsarClientException { protected CompletableFuture> internalBatchReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - try { - lock.writeLock().lock(); - if (pendingBatchReceives == null) { - pendingBatchReceives = Queues.newConcurrentLinkedQueue(); - } - if (hasEnoughMessagesForBatchReceive()) { - MessagesImpl messages = getNewMessagesImpl(); - Message msgPeeked = incomingMessages.peek(); - while (msgPeeked != null && messages.canAdd(msgPeeked)) { - Message msg = incomingMessages.poll(); - if (msg != null) { - decreaseIncomingMessageSize(msg); - Message interceptMsg = beforeConsume(msg); - messages.add(interceptMsg); - } - msgPeeked = incomingMessages.peek(); + if (pendingBatchReceives == null) { + pendingBatchReceives = Queues.newConcurrentLinkedQueue(); + } + if (hasEnoughMessagesForBatchReceive()) { + MessagesImpl messages = getNewMessagesImpl(); + Message msg = null; + do { + msg = incomingMessages.poll(); + if (msg != null) { + decreaseIncomingMessageSize(msg); + Message interceptMsg = beforeConsume(msg); + messages.add(interceptMsg); } - result.complete(messages); - } else { - OpBatchReceive opBatchReceive = OpBatchReceive.of(result); - pendingBatchReceives.add(opBatchReceive); - cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); - } - resumeReceivingFromPausedConsumersIfNeeded(); - } finally { - lock.writeLock().unlock(); + } while (msg != null && messages.canAdd()); + result.complete(messages); + } else { + OpBatchReceive opBatchReceive = OpBatchReceive.of(result); + pendingBatchReceives.add(opBatchReceive); + cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); } + resumeReceivingFromPausedConsumersIfNeeded(); + return result; } @@ -595,18 +587,14 @@ private ConsumerConfigurationData getInternalConsumerConfig() { @Override public void redeliverUnacknowledgedMessages() { - lock.writeLock().lock(); - try { - consumers.values().stream().forEach(consumer -> { - consumer.redeliverUnacknowledgedMessages(); - consumer.unAckedChunkedMessageIdSequenceMap.clear(); - }); - incomingMessages.clear(); - resetIncomingMessageSize(); - unAckedMessageTracker.clear(); - } finally { - lock.writeLock().unlock(); - } + consumers.values().stream().forEach(consumer -> { + consumer.redeliverUnacknowledgedMessages(); + consumer.unAckedChunkedMessageIdSequenceMap.clear(); + }); + incomingMessages.clear(); + resetIncomingMessageSize(); + unAckedMessageTracker.clear(); + resumeReceivingFromPausedConsumersIfNeeded(); } From bce28556abf95dab4b1e9137addebd94e5391435 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sun, 24 Jan 2021 17:01:03 +0800 Subject: [PATCH 2/3] move lock back for batchReceive --- .../pulsar/client/impl/ConsumerBase.java | 20 ++++++--- .../pulsar/client/impl/ConsumerImpl.java | 45 +++++++++++-------- .../pulsar/client/impl/MessagesImpl.java | 12 ----- 3 files changed, 40 insertions(+), 37 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java index 33b721c27683c..28c248fe674ac 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBase.java @@ -37,6 +37,8 @@ import java.util.concurrent.atomic.AtomicLongFieldUpdater; import io.netty.util.Timeout; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.ConsumerEventListener; @@ -82,6 +84,7 @@ public abstract class ConsumerBase extends HandlerState implements Consumer conf, int receiverQueueSize, ExecutorProvider executorProvider, @@ -731,7 +734,12 @@ protected void notifyPendingBatchReceivedCallBack() { return; } - notifyPendingBatchReceivedCallBack(opBatchReceive); + try { + reentrantLock.lock(); + notifyPendingBatchReceivedCallBack(opBatchReceive); + } finally { + reentrantLock.unlock(); + } } private OpBatchReceive peekNextBatchReceive() { @@ -773,16 +781,16 @@ private OpBatchReceive pollNextBatchReceive() { protected final void notifyPendingBatchReceivedCallBack(OpBatchReceive opBatchReceive) { MessagesImpl messages = getNewMessagesImpl(); - - Message msg = null; - do { - msg = incomingMessages.poll(); + Message msgPeeked = incomingMessages.peek(); + while (msgPeeked != null && messages.canAdd(msgPeeked)) { + Message msg = incomingMessages.poll(); if (msg != null) { messageProcessed(msg); Message interceptMsg = beforeConsume(msg); messages.add(interceptMsg); } - } while (msg != null && messages.canAdd()); + msgPeeked = incomingMessages.peek(); + } completePendingBatchReceive(opBatchReceive.future, messages); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index fb9bc624bee9c..d57fc39a7d97b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -126,6 +126,8 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final int receiverQueueRefillThreshold; + private final ReadWriteLock lock = new ReentrantReadWriteLock(); + private final UnAckedMessageTracker unAckedMessageTracker; private final AcknowledgmentsGroupingTracker acknowledgmentsGroupingTracker; private final NegativeAcksTracker negativeAcksTracker; @@ -465,26 +467,31 @@ protected Messages internalBatchReceive() throws PulsarClientException { protected CompletableFuture> internalBatchReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - if (pendingBatchReceives == null) { - pendingBatchReceives = Queues.newConcurrentLinkedQueue(); - } - if (hasEnoughMessagesForBatchReceive()) { - MessagesImpl messages = getNewMessagesImpl(); - Message msg = null; - do { - msg = incomingMessages.poll(); - if (msg != null) { - messageProcessed(msg); - Message interceptMsg = beforeConsume(msg); - messages.add(interceptMsg); + try { + lock.writeLock().lock(); + if (pendingBatchReceives == null) { + pendingBatchReceives = Queues.newConcurrentLinkedQueue(); + } + if (hasEnoughMessagesForBatchReceive()) { + MessagesImpl messages = getNewMessagesImpl(); + Message msgPeeked = incomingMessages.peek(); + while (msgPeeked != null && messages.canAdd(msgPeeked)) { + Message msg = incomingMessages.poll(); + if (msg != null) { + messageProcessed(msg); + Message interceptMsg = beforeConsume(msg); + messages.add(interceptMsg); + } + msgPeeked = incomingMessages.peek(); } - } while (msg != null && messages.canAdd()); - - result.complete(messages); - } else { - OpBatchReceive opBatchReceive = OpBatchReceive.of(result); - pendingBatchReceives.add(opBatchReceive); - cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); + result.complete(messages); + } else { + OpBatchReceive opBatchReceive = OpBatchReceive.of(result); + pendingBatchReceives.add(opBatchReceive); + cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); + } + } finally { + lock.writeLock().unlock(); } return result; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java index f2705958784ef..c56694e2759ea 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java @@ -56,18 +56,6 @@ protected boolean canAdd(Message message) { return true; } - protected boolean canAdd() { - if (maxNumberOfMessages > 0 && currentNumberOfMessages + 1 > maxNumberOfMessages) { - return false; - } - - if (maxSizeOfMessages > 0 && currentSizeOfMessages + 1 > maxSizeOfMessages) { - return false; - } - - return true; - } - protected void add(Message message) { if (message == null) { return; From 6f211738454362fcc3cfaf12ba363f43bec5a2d8 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Sun, 24 Jan 2021 21:51:29 +0800 Subject: [PATCH 3/3] move lock back for batchReceive --- .../client/impl/MultiTopicsConsumerImpl.java | 47 +++++++++++-------- 1 file changed, 28 insertions(+), 19 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java index 4c820dd680e31..52634019d4e50 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java @@ -62,6 +62,8 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.ReadWriteLock; +import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -95,6 +97,7 @@ public class MultiTopicsConsumerImpl extends ConsumerBase { private volatile Timeout partitionsAutoUpdateTimeout = null; TopicsPartitionChangedListener topicsPartitionChangedListener; CompletableFuture partitionsAutoUpdateFuture = null; + private final ReadWriteLock lock = new ReentrantReadWriteLock(); private final ConsumerStatsRecorder stats; private final UnAckedMessageTracker unAckedMessageTracker; @@ -356,27 +359,33 @@ protected Messages internalBatchReceive() throws PulsarClientException { protected CompletableFuture> internalBatchReceiveAsync() { CompletableFutureCancellationHandler cancellationHandler = new CompletableFutureCancellationHandler(); CompletableFuture> result = cancellationHandler.createFuture(); - if (pendingBatchReceives == null) { - pendingBatchReceives = Queues.newConcurrentLinkedQueue(); - } - if (hasEnoughMessagesForBatchReceive()) { - MessagesImpl messages = getNewMessagesImpl(); - Message msg = null; - do { - msg = incomingMessages.poll(); - if (msg != null) { - decreaseIncomingMessageSize(msg); - Message interceptMsg = beforeConsume(msg); - messages.add(interceptMsg); + try { + lock.writeLock().lock(); + if (pendingBatchReceives == null) { + pendingBatchReceives = Queues.newConcurrentLinkedQueue(); + } + if (hasEnoughMessagesForBatchReceive()) { + MessagesImpl messages = getNewMessagesImpl(); + Message msgPeeked = incomingMessages.peek(); + while (msgPeeked != null && messages.canAdd(msgPeeked)) { + Message msg = incomingMessages.poll(); + if (msg != null) { + decreaseIncomingMessageSize(msg); + Message interceptMsg = beforeConsume(msg); + messages.add(interceptMsg); + } + msgPeeked = incomingMessages.peek(); } - } while (msg != null && messages.canAdd()); - result.complete(messages); - } else { - OpBatchReceive opBatchReceive = OpBatchReceive.of(result); - pendingBatchReceives.add(opBatchReceive); - cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); + result.complete(messages); + } else { + OpBatchReceive opBatchReceive = OpBatchReceive.of(result); + pendingBatchReceives.add(opBatchReceive); + cancellationHandler.setCancelAction(() -> pendingBatchReceives.remove(opBatchReceive)); + } + resumeReceivingFromPausedConsumersIfNeeded(); + } finally { + lock.writeLock().unlock(); } - resumeReceivingFromPausedConsumersIfNeeded(); return result; }