diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SharedConsumerAssignor.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SharedConsumerAssignor.java index 4814d0e22150e..ddefd5abc81b5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SharedConsumerAssignor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/SharedConsumerAssignor.java @@ -123,16 +123,44 @@ private Consumer getConsumer(final int numConsumers) { return null; } + /** + * Remove every uuid mapping that points to the given consumer. Must be called when a consumer is removed from the + * dispatcher: a mapping to a disconnected consumer has no permits left, so the remaining chunks of that uuid stay + * unassignable and stall the subscription. Also bounds the map when the last chunk is never published. + */ + public void removeConsumer(final Consumer consumer) { + uuidToConsumer.values().removeIf(cachedConsumer -> cachedConsumer == consumer); + } + + /** Drop all uuid mappings, e.g. after all consumers have been removed from the dispatcher. */ + public void clear() { + uuidToConsumer.clear(); + } + private Consumer getConsumerForUuid(final MessageMetadata metadata, final Consumer defaultConsumer) { final String uuid = metadata.getUuid(); Consumer consumer = uuidToConsumer.get(uuid); if (consumer == null) { - if (metadata.getChunkId() != 0) { - // Not the first chunk, skip it - return null; - } consumer = defaultConsumer; - uuidToConsumer.put(uuid, consumer); + if (metadata.getChunkId() == 0) { + uuidToConsumer.put(uuid, consumer); + } else { + // Orphan chunk: only chunk 0 creates the mapping, and the cursor always re-reads chunk 0 before + // this one while it is still unacked, so a missing mapping means chunk 0 is acknowledged and gone. + // Never replay it: the dispatcher drains the replay queue before every normal read, so an entry that + // can never be assigned stalls the whole subscription. Dispatch it and let the client discard it. + if (subscription != null) { + log.warn() + .attr("topic", subscription.getTopic().getName()) + .attr("subscription", subscription.getName()) + .attr("uuid", uuid) + .attr("chunkId", metadata.getChunkId()) + .attr("numChunks", metadata.getNumChunksFromMsg()) + .attr("consumer", defaultConsumer) + .log("Dispatching orphan chunk whose first chunk is gone. The client should discard and" + + " acknowledge it"); + } + } } final int permits = consumerToPermits.computeIfAbsent(consumer, Consumer::getAvailablePermits); if (permits <= 0) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index b659f6e2200d8..5a1ea4a559b5f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -241,6 +241,8 @@ protected boolean isConsumersExceededOnSubscription() { public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { // decrement unack-message count for removed consumer addUnAckedMessages(-consumer.getUnackedMessages()); + // Drop this consumer's chunk uuid mappings, otherwise its pending chunks stay unassignable forever. + assignor.removeConsumer(consumer); if (consumerSet.removeAll(consumer) == 1) { consumerList.remove(consumer); log.info() @@ -291,6 +293,7 @@ protected synchronized void clearComponentsAfterRemovedAllConsumers() { redeliveryMessages.clear(); redeliveryTracker.clear(); + assignor.clear(); if (closeFuture != null) { log.info("All consumers removed. Subscription is disconnected"); closeFuture.complete(null); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java index 3de50042b592d..0a2f742729269 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumersClassic.java @@ -231,6 +231,8 @@ protected boolean isConsumersExceededOnSubscription() { public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { // decrement unack-message count for removed consumer addUnAckedMessages(-consumer.getUnackedMessages()); + // Drop this consumer's chunk uuid mappings, otherwise its pending chunks stay unassignable forever. + assignor.removeConsumer(consumer); if (consumerSet.removeAll(consumer) == 1) { consumerList.remove(consumer); log.info() @@ -271,6 +273,7 @@ private synchronized void clearComponentsAfterRemovedAllConsumers() { redeliveryMessages.clear(); redeliveryTracker.clear(); + assignor.clear(); if (closeFuture != null) { log.info("All consumers removed. Subscription is disconnected"); closeFuture.complete(null); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedConsumerAssignorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedConsumerAssignorTest.java index d9e675fa7e543..a77c87536c942 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedConsumerAssignorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/SharedConsumerAssignorTest.java @@ -153,6 +153,88 @@ public void testMultiConsumerWithSmallPermits() { assertTrue(assignor.getUuidToConsumer().isEmpty()); } + /** + * An orphan chunk is a chunk whose uuid is not tracked and whose chunk id is not 0, e.g. entry 5 (A-1-2-3) when + * the earlier chunks of "A-1" were acknowledged before the mapping existed. That happens after a topic unload, a + * broker restart, or when the consumer that owned the uuid disconnected. + *
+ * Such an entry can never become assignable, so handing it to the unassigned message processor (the replay queue)
+ * stalls the subscription forever: {@code PersistentDispatcherMultipleConsumers#readMoreEntries} drains the replay
+ * queue before every normal read, so the dispatcher keeps replaying the same entry, never reads anything new and
+ * the mark-delete position never moves. It must be dispatched instead, so that the client can discard and
+ * acknowledge the incomplete chunked message.
+ */
+ @Test
+ public void testOrphanChunkIsDispatchedInsteadOfReplayedForever() {
+ final Consumer consumer = new Consumer("A", 100);
+ roundRobinConsumerSelector.addConsumers(consumer);
+ assertTrue(assignor.getUuidToConsumer().isEmpty());
+
+ // The last chunk of "A-1" is read while the uuid was never cached
+ Map