diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java index 7046ba481937a..0280ceda725c4 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java @@ -194,4 +194,11 @@ public synchronized Throwable fillInStackTrace() { // Disable stack traces to be filled in return null; } + + public static class OffloadReadHandleClosedException extends ManagedLedgerException { + + public OffloadReadHandleClosedException() { + super("Offload read handle already closed"); + } + } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 9a92d35c057af..6f933962d4651 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -2469,7 +2469,7 @@ void internalTrimLedgers(boolean isTruncate, CompletableFuture promise) { break; } // if truncate, all ledgers besides currentLedger are going to be deleted - if (isTruncate){ + if (isTruncate) { if (log.isDebugEnabled()) { log.debug("[{}] Ledger {} will be truncated with ts {}", name, ls.getLedgerId(), ls.getTimestamp()); @@ -2497,11 +2497,14 @@ void internalTrimLedgers(boolean isTruncate, CompletableFuture promise) { } ledgersToDelete.add(ls); } else { - // once retention constraint has been met, skip check - if (log.isDebugEnabled()) { - log.debug("[{}] Ledger {} not deleted. Neither expired nor over-quota", name, ls.getLedgerId()); + if (ls.getLedgerId() < getTheSlowestNonDurationReadPosition().getLedgerId()) { + // once retention constraint has been met, skip check + if (log.isDebugEnabled()) { + log.debug("[{}] Ledger {} not deleted. Neither expired nor over-quota", name, + ls.getLedgerId()); + } + invalidateReadHandle(ls.getLedgerId()); } - invalidateReadHandle(ls.getLedgerId()); } } @@ -4149,4 +4152,17 @@ public void checkInactiveLedgerAndRollOver() { } } + public Position getTheSlowestNonDurationReadPosition() { + PositionImpl theSlowestNonDurableReadPosition = PositionImpl.LATEST; + for (ManagedCursor cursor : cursors) { + if (cursor instanceof NonDurableCursorImpl) { + PositionImpl readPosition = (PositionImpl) cursor.getReadPosition(); + if (readPosition.compareTo(theSlowestNonDurableReadPosition) < 0) { + theSlowestNonDurableReadPosition = readPosition; + } + } + } + return theSlowestNonDurableReadPosition; + } + } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index 908d381a8387f..753c34f38ef2b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -18,6 +18,7 @@ */ package org.apache.bookkeeper.mledger.impl; +import static java.nio.charset.StandardCharsets.UTF_8; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyString; @@ -3580,6 +3581,31 @@ public void testOffloadTaskCancelled() throws Exception { }); } + @Test + public void testGetTheSlowestNonDurationReadPosition() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("test_", + new ManagedLedgerConfig().setMaxEntriesPerLedger(1).setRetentionTime(-1, TimeUnit.SECONDS) + .setRetentionSizeInMB(-1)); + ledger.openCursor("c1"); + + List positions = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + positions.add(ledger.addEntry(("entry-" + i).getBytes(UTF_8))); + } + + Assert.assertEquals(ledger.getTheSlowestNonDurationReadPosition(), PositionImpl.LATEST); + + ManagedCursor nonDurableCursor = ledger.newNonDurableCursor(PositionImpl.EARLIEST); + + Assert.assertEquals(ledger.getTheSlowestNonDurationReadPosition(), positions.get(0)); + + ledger.deleteCursor(nonDurableCursor.getName()); + + Assert.assertEquals(ledger.getTheSlowestNonDurationReadPosition(), PositionImpl.LATEST); + + ledger.close(); + } + @Test public void testGetLedgerMetadata() throws Exception { ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) factory.open("testGetLedgerMetadata"); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorTest.java index 697cf28c172d3..437ca9937afc3 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/NonDurableCursorTest.java @@ -31,6 +31,7 @@ import com.google.common.collect.Lists; import java.nio.charset.Charset; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; @@ -54,6 +55,7 @@ import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.testng.Assert; import org.testng.annotations.Test; public class NonDurableCursorTest extends MockedBookKeeperTestCase { @@ -735,6 +737,64 @@ public void testBacklogStatsWhenDroppingData() throws Exception { ledger.close(); } + @Test + public void testInvalidateReadHandleWithSlowNonDurableCursor() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testInvalidateReadHandleWithSlowNonDurableCursor", + new ManagedLedgerConfig().setMaxEntriesPerLedger(1).setRetentionTime(-1, TimeUnit.SECONDS) + .setRetentionSizeInMB(-1)); + ManagedCursor c1 = ledger.openCursor("c1"); + ManagedCursor nonDurableCursor = ledger.newNonDurableCursor(PositionImpl.EARLIEST); + + List positions = new ArrayList<>(); + for (int i = 0; i < 10; i++) { + positions.add(ledger.addEntry(("entry-" + i).getBytes(UTF_8))); + } + + CountDownLatch latch = new CountDownLatch(10); + for (int i = 0; i < 10; i++) { + ledger.asyncReadEntry((PositionImpl) positions.get(i), new AsyncCallbacks.ReadEntryCallback() { + @Override + public void readEntryComplete(Entry entry, Object ctx) { + latch.countDown(); + } + + @Override + public void readEntryFailed(ManagedLedgerException exception, Object ctx) { + latch.countDown(); + } + }, null); + } + + latch.await(); + + c1.markDelete(positions.get(4)); + + CompletableFuture promise = new CompletableFuture<>(); + ledger.internalTrimConsumedLedgers(promise); + promise.join(); + + Assert.assertTrue(ledger.ledgerCache.containsKey(positions.get(0).getLedgerId())); + Assert.assertTrue(ledger.ledgerCache.containsKey(positions.get(1).getLedgerId())); + Assert.assertTrue(ledger.ledgerCache.containsKey(positions.get(2).getLedgerId())); + Assert.assertTrue(ledger.ledgerCache.containsKey(positions.get(3).getLedgerId())); + Assert.assertTrue(ledger.ledgerCache.containsKey(positions.get(4).getLedgerId())); + + promise = new CompletableFuture<>(); + + nonDurableCursor.markDelete(positions.get(3)); + + ledger.internalTrimConsumedLedgers(promise); + promise.join(); + + Assert.assertFalse(ledger.ledgerCache.containsKey(positions.get(0).getLedgerId())); + Assert.assertFalse(ledger.ledgerCache.containsKey(positions.get(1).getLedgerId())); + Assert.assertFalse(ledger.ledgerCache.containsKey(positions.get(2).getLedgerId())); + Assert.assertFalse(ledger.ledgerCache.containsKey(positions.get(3).getLedgerId())); + Assert.assertTrue(ledger.ledgerCache.containsKey(positions.get(4).getLedgerId())); + + ledger.close(); + } + @Test(expectedExceptions = NullPointerException.class) void testCursorWithNameIsNotNull() throws Exception { final String p1CursorName = "entry-1"; 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 ee11d8961b8b9..bfff0c67c73fd 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 @@ -657,7 +657,8 @@ public synchronized void readEntriesFailed(ManagedLedgerException exception, Obj // Notify the consumer only if all the messages were already acknowledged consumerList.forEach(Consumer::reachedEndOfTopic); } - } else if (exception.getCause() instanceof TransactionBufferException.TransactionNotSealedException) { + } else if (exception.getCause() instanceof TransactionBufferException.TransactionNotSealedException + || exception.getCause() instanceof ManagedLedgerException.OffloadReadHandleClosedException) { waitTimeMillis = 1; if (log.isDebugEnabled()) { log.debug("[{}] Error reading transaction entries : {}, Read Type {} - Retrying to read in {} seconds", diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 79f397af88322..1f2636fbe4ca6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -497,7 +497,8 @@ private synchronized void internalReadEntriesFailed(ManagedLedgerException excep // Notify the consumer only if all the messages were already acknowledged consumers.forEach(Consumer::reachedEndOfTopic); } - } else if (exception.getCause() instanceof TransactionBufferException.TransactionNotSealedException) { + } else if (exception.getCause() instanceof TransactionBufferException.TransactionNotSealedException + || exception.getCause() instanceof ManagedLedgerException.OffloadReadHandleClosedException) { waitTimeMillis = 1; if (log.isDebugEnabled()) { log.debug("[{}] Error reading transaction entries : {}, - Retrying to read in {} seconds", name, diff --git a/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImpl.java b/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImpl.java index fae534454e0d7..c5ca5e1aaf3a6 100644 --- a/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImpl.java +++ b/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImpl.java @@ -35,6 +35,7 @@ import org.apache.bookkeeper.client.api.ReadHandle; import org.apache.bookkeeper.client.impl.LedgerEntriesImpl; import org.apache.bookkeeper.client.impl.LedgerEntryImpl; +import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.offload.jcloud.BackedInputStream; import org.apache.bookkeeper.mledger.offload.jcloud.OffloadIndexBlock; import org.apache.bookkeeper.mledger.offload.jcloud.OffloadIndexBlockBuilder; @@ -103,14 +104,15 @@ public CompletableFuture readAsync(long firstEntry, long lastEntr log.debug("Ledger {}: reading {} - {}", getId(), firstEntry, lastEntry); CompletableFuture promise = new CompletableFuture<>(); executor.submit(() -> { + if (state == State.Closed) { + log.warn("Reading a closed read handler. Ledger ID: {}, Read range: {}-{}", + ledgerId, firstEntry, lastEntry); + promise.completeExceptionally(new ManagedLedgerException.OffloadReadHandleClosedException()); + return; + } List entries = new ArrayList(); boolean seeked = false; try { - if (state == State.Closed) { - log.warn("Reading a closed read handler. Ledger ID: {}, Read range: {}-{}", - ledgerId, firstEntry, lastEntry); - throw new BKException.BKUnexpectedConditionException(); - } if (firstEntry > lastEntry || firstEntry < 0 || lastEntry > getLastAddConfirmed()) { diff --git a/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImplV2.java b/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImplV2.java index a6f66132c4af1..a99f0abd24892 100644 --- a/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImplV2.java +++ b/tiered-storage/jcloud/src/main/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreBackedReadHandleImplV2.java @@ -38,6 +38,7 @@ import org.apache.bookkeeper.client.api.ReadHandle; import org.apache.bookkeeper.client.impl.LedgerEntriesImpl; import org.apache.bookkeeper.client.impl.LedgerEntryImpl; +import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.offload.jcloud.BackedInputStream; import org.apache.bookkeeper.mledger.offload.jcloud.OffloadIndexBlockV2; import org.apache.bookkeeper.mledger.offload.jcloud.OffloadIndexBlockV2Builder; @@ -57,6 +58,13 @@ public class BlobStoreBackedReadHandleImplV2 implements ReadHandle { private final List dataStreams; private final ExecutorService executor; + private State state = null; + + enum State { + Opened, + Closed + } + static class GroupedReader { @Override public String toString() { @@ -97,6 +105,7 @@ private BlobStoreBackedReadHandleImplV2(long ledgerId, List dataStreams.add(new DataInputStream(inputStream)); } this.executor = executor; + this.state = State.Opened; } @Override @@ -121,6 +130,7 @@ public CompletableFuture closeAsync() { for (DataInputStream dataStream : dataStreams) { dataStream.close(); } + state = State.Closed; promise.complete(null); } catch (IOException t) { promise.completeExceptionally(t); @@ -133,13 +143,20 @@ public CompletableFuture closeAsync() { public CompletableFuture readAsync(long firstEntry, long lastEntry) { log.debug("Ledger {}: reading {} - {}", getId(), firstEntry, lastEntry); CompletableFuture promise = new CompletableFuture<>(); - if (firstEntry > lastEntry - || firstEntry < 0 - || lastEntry > getLastAddConfirmed()) { - promise.completeExceptionally(new IllegalArgumentException()); - return promise; - } executor.submit(() -> { + if (state == State.Closed) { + log.warn("Reading a closed read handler. Ledger ID: {}, Read range: {}-{}", + ledgerId, firstEntry, lastEntry); + promise.completeExceptionally(new ManagedLedgerException.OffloadReadHandleClosedException()); + return; + } + + if (firstEntry > lastEntry + || firstEntry < 0 + || lastEntry > getLastAddConfirmed()) { + promise.completeExceptionally(new BKException.BKIncorrectParameterException()); + return; + } List entries = new ArrayList(); List groupedReaders = null; try { diff --git a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreManagedLedgerOffloaderTest.java b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreManagedLedgerOffloaderTest.java index ab979f8a5a142..b78666ad6454f 100644 --- a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreManagedLedgerOffloaderTest.java +++ b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/BlobStoreManagedLedgerOffloaderTest.java @@ -41,6 +41,7 @@ import org.apache.bookkeeper.client.api.LedgerEntry; import org.apache.bookkeeper.client.api.ReadHandle; import org.apache.bookkeeper.mledger.LedgerOffloader; +import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.offload.jcloud.provider.JCloudBlobStoreProvider; import org.apache.bookkeeper.mledger.offload.jcloud.provider.TieredStorageConfiguration; import org.jclouds.blobstore.BlobStore; @@ -512,7 +513,7 @@ public void testReadWithAClosedLedgerHandler() throws Exception { try { toTest.readAsync(0, lac).get(); } catch (Exception e) { - if (e.getCause() instanceof BKException.BKUnexpectedConditionException) { + if (e.getCause() instanceof ManagedLedgerException.OffloadReadHandleClosedException) { // expected exception return; }