diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawBatchConverter.java b/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawBatchConverter.java index e252426acbbfa..8c21a737e65be 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawBatchConverter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/client/impl/RawBatchConverter.java @@ -91,8 +91,7 @@ public static List> extractIdsAndKeys(RawMessage * Take a batched message and a filter, and returns a message with the only the sub-messages * which match the filter. Returns an empty optional if no messages match. * - * This takes ownership of the passes in message, and if the returned optional is not empty, - * the ownership of that message is returned also. + * NOTE: this message does not alter the reference count of the RawMessage argument. */ public static Optional rebatchMessage(RawMessage msg, BiPredicate filter) @@ -161,9 +160,9 @@ public static Optional rebatchMessage(RawMessage msg, return Optional.empty(); } } finally { + uncompressedPayload.release(); batchBuffer.release(); metadata.recycle(); - msg.close(); } } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java index b1378b648bf81..22efe8e5c5719 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java @@ -164,12 +164,19 @@ private static CompletableFuture readOneMessageId(LedgerHandle lh if (rc != BKException.Code.OK) { promise.completeExceptionally(BKException.create(rc)); } else { - try (RawMessage m = RawMessageImpl.deserializeFrom( - seq.nextElement().getEntryBuffer())) { - promise.complete(m.getMessageIdData()); - } catch (NoSuchElementException e) { - log.error("No such entry {} in ledger {}", entryId, lh.getId()); - promise.completeExceptionally(e); + // Need to release buffers for all entries in the sequence + if (seq.hasMoreElements()) { + LedgerEntry entry = seq.nextElement(); + try (RawMessage m = RawMessageImpl.deserializeFrom(entry.getEntryBuffer())) { + entry.getEntryBuffer().release(); + while (seq.hasMoreElements()) { + seq.nextElement().getEntryBuffer().release(); + } + promise.complete(m.getMessageIdData()); + } + } else { + promise.completeExceptionally(new NoSuchElementException( + String.format("No such entry %d in ledger %d", entryId, lh.getId()))); } } }, null); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java index 95f6f1ad7f5ad..a275bb5fb107e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java @@ -212,77 +212,88 @@ private CompletableFuture phaseTwoSeekThenLoop(RawReader reader, MessageId private void phaseTwoLoop(RawReader reader, MessageId to, Map latestForKey, LedgerHandle lh, Semaphore outstanding, CompletableFuture promise) { + if (promise.isDone()) { + return; + } reader.readNextAsync().whenCompleteAsync( (m, exception) -> { if (exception != null) { promise.completeExceptionally(exception); return; } else if (promise.isDone()) { + m.close(); return; } - MessageId id = m.getMessageId(); - Optional messageToAdd = Optional.empty(); - if (RawBatchConverter.isReadableBatch(m)) { - try { - messageToAdd = RawBatchConverter.rebatchMessage( - m, (key, subid) -> latestForKey.get(key).equals(subid)); - } catch (IOException ioe) { - log.info("Error decoding batch for message {}. Whole batch will be included in output", - id, ioe); - messageToAdd = Optional.of(m); - } - } else { - Pair keyAndSize = extractKeyAndSize(m); - MessageId msg; - if (keyAndSize == null) { // pass through messages without a key - messageToAdd = Optional.of(m); - } else if ((msg = latestForKey.get(keyAndSize.getLeft())) != null - && msg.equals(id)) { // consider message only if present into latestForKey map - if (keyAndSize.getRight() <= 0) { - promise.completeExceptionally(new IllegalArgumentException( - "Compaction phase found empty record from sorted key-map")); + try { + MessageId id = m.getMessageId(); + Optional messageToAdd = Optional.empty(); + if (RawBatchConverter.isReadableBatch(m)) { + try { + messageToAdd = RawBatchConverter.rebatchMessage( + m, (key, subid) -> latestForKey.get(key).equals(subid)); + } catch (IOException ioe) { + log.info("Error decoding batch for message {}. Whole batch will be included in output", + id, ioe); + messageToAdd = Optional.of(m); } - messageToAdd = Optional.of(m); } else { - m.close(); + Pair keyAndSize = extractKeyAndSize(m); + MessageId msg; + if (keyAndSize == null) { // pass through messages without a key + messageToAdd = Optional.of(m); + } else if ((msg = latestForKey.get(keyAndSize.getLeft())) != null + && msg.equals(id)) { // consider message only if present into latestForKey map + if (keyAndSize.getRight() <= 0) { + promise.completeExceptionally(new IllegalArgumentException( + "Compaction phase found empty record from sorted key-map")); + } + messageToAdd = Optional.of(m); + } } - } - if (messageToAdd.isPresent()) { - try { - outstanding.acquire(); - CompletableFuture addFuture = addToCompactedLedger(lh, messageToAdd.get()) - .whenComplete((res, exception2) -> { - outstanding.release(); - if (exception2 != null) { - promise.completeExceptionally(exception2); + if (messageToAdd.isPresent()) { + RawMessage message = messageToAdd.get(); + try { + outstanding.acquire(); + CompletableFuture addFuture = addToCompactedLedger(lh, message) + .whenComplete((res, exception2) -> { + outstanding.release(); + if (exception2 != null) { + promise.completeExceptionally(exception2); + } + }); + if (to.equals(id)) { + addFuture.whenComplete((res, exception2) -> { + if (exception2 == null) { + promise.complete(null); } }); - if (to.equals(id)) { - addFuture.whenComplete((res, exception2) -> { - if (exception2 == null) { - promise.complete(null); - } - }); + } + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + promise.completeExceptionally(ie); + } finally { + if (message != m) { + message.close(); + } } - } catch (InterruptedException ie) { - Thread.currentThread().interrupt(); - promise.completeExceptionally(ie); - } - } else if (to.equals(id)) { - // Reached to last-id and phase-one found it deleted-message while iterating on ledger so, not - // present under latestForKey. Complete the compaction. - try { - // make sure all inflight writes have finished - outstanding.acquire(MAX_OUTSTANDING); - promise.complete(null); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - promise.completeExceptionally(e); + } else if (to.equals(id)) { + // Reached to last-id and phase-one found it deleted-message while iterating on ledger so, + // not present under latestForKey. Complete the compaction. + try { + // make sure all inflight writes have finished + outstanding.acquire(MAX_OUTSTANDING); + promise.complete(null); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + promise.completeExceptionally(e); + } + return; } - return; + phaseTwoLoop(reader, to, latestForKey, lh, outstanding, promise); + } finally { + m.close(); } - phaseTwoLoop(reader, to, latestForKey, lh, outstanding, promise); }, scheduler); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/RawReaderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/RawReaderTest.java index b0c7cd1830055..5ae46185c4cb0 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/RawReaderTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/RawReaderTest.java @@ -319,13 +319,13 @@ public void testBatchingRebatch() throws Exception { } RawReader reader = RawReader.create(pulsarClient, topic, subscription).get(); - try { - RawMessage m1 = reader.readNextAsync().get(); + try (RawMessage m1 = reader.readNextAsync().get()) { RawMessage m2 = RawBatchConverter.rebatchMessage(m1, (key, id) -> key.equals("key2")).get(); List> idsAndKeys = RawBatchConverter.extractIdsAndKeys(m2); Assert.assertEquals(idsAndKeys.size(), 1); Assert.assertEquals(idsAndKeys.get(0).getRight(), "key2"); m2.close(); + Assert.assertEquals(m1.getHeadersAndPayload().refCnt(), 1); } finally { reader.closeAsync().get(); }