Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down Expand Up @@ -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());
}
}

Expand Down Expand Up @@ -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;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Position> 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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 {
Expand Down Expand Up @@ -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<Position> 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<Void> 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";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -103,14 +104,15 @@ public CompletableFuture<LedgerEntries> readAsync(long firstEntry, long lastEntr
log.debug("Ledger {}: reading {} - {}", getId(), firstEntry, lastEntry);
CompletableFuture<LedgerEntries> 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<LedgerEntry> entries = new ArrayList<LedgerEntry>();
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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -57,6 +58,13 @@ public class BlobStoreBackedReadHandleImplV2 implements ReadHandle {
private final List<DataInputStream> dataStreams;
private final ExecutorService executor;

private State state = null;

enum State {
Opened,
Closed
}

static class GroupedReader {
@Override
public String toString() {
Expand Down Expand Up @@ -97,6 +105,7 @@ private BlobStoreBackedReadHandleImplV2(long ledgerId, List<OffloadIndexBlockV2>
dataStreams.add(new DataInputStream(inputStream));
}
this.executor = executor;
this.state = State.Opened;
}

@Override
Expand All @@ -121,6 +130,7 @@ public CompletableFuture<Void> closeAsync() {
for (DataInputStream dataStream : dataStreams) {
dataStream.close();
}
state = State.Closed;
promise.complete(null);
} catch (IOException t) {
promise.completeExceptionally(t);
Expand All @@ -133,13 +143,20 @@ public CompletableFuture<Void> closeAsync() {
public CompletableFuture<LedgerEntries> readAsync(long firstEntry, long lastEntry) {
log.debug("Ledger {}: reading {} - {}", getId(), firstEntry, lastEntry);
CompletableFuture<LedgerEntries> 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<LedgerEntry> entries = new ArrayList<LedgerEntry>();
List<GroupedReader> groupedReaders = null;
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down