Skip to content
Open
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 @@ -620,6 +620,11 @@ void asyncOpenCursor(String name, InitialPosition initialPosition, Map<String, L
void asyncSetProperties(Map<String, String> properties, AsyncCallbacks.UpdatePropertiesCallback callback,
Object ctx);

/**
* Reset cursor if ledger consumed completely, before trim consumed ledgers in background.
*/
boolean maybeUpdateCursorBeforeTrimmingConsumedLedger();

/**
* Trim consumed ledgers in background.
* @param promise
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2296,7 +2296,8 @@ public void addWaitingEntryCallBack(WaitingEntryCallBack cb) {
this.waitingEntryCallBacks.add(cb);
}

public void maybeUpdateCursorBeforeTrimmingConsumedLedger() {
public boolean maybeUpdateCursorBeforeTrimmingConsumedLedger() {
boolean maybeUpdatedCursor = false;
for (ManagedCursor cursor : cursors) {
PositionImpl lastAckedPosition = (PositionImpl) cursor.getMarkDeletedPosition();
LedgerInfo currPointedLedger = ledgers.get(lastAckedPosition.getLedgerId());
Expand All @@ -2313,20 +2314,22 @@ public void maybeUpdateCursorBeforeTrimmingConsumedLedger() {
log.debug("No need to reset cursor: {}, current ledger is the last ledger.", cursor);
}
} else {
log.warn("Cursor: {} does not exist in the managed-ledger.", cursor);
log.debug("No need to reset cursor: {}, current ledger maybe has removed from managed-ledger.", cursor);
}

if (!lastAckedPosition.equals((PositionImpl) cursor.getMarkDeletedPosition())) {
try {
log.info("Reset cursor:{} to {} since ledger consumed completely", cursor, lastAckedPosition);
updateCursor((ManagedCursorImpl) cursor, lastAckedPosition);
maybeUpdatedCursor = true;
} catch (Exception e) {
log.warn("Failed to reset cursor: {} from {} to {}. Trimming thread will retry next time.",
cursor, cursor.getMarkDeletedPosition(), lastAckedPosition);
log.warn("Caused by", e);
}
}
}
return maybeUpdatedCursor;
}

private void trimConsumedLedgersInBackground() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1782,7 +1782,11 @@ private void checkConsumedLedgers() {
if (t instanceof PersistentTopic) {
Optional.ofNullable(((PersistentTopic) t).getManagedLedger()).ifPresent(
managedLedger -> {
managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE);
// After update cursor, trimConsumedLedgersInBackground will be invoked,
// avoid invoke repeatedly here.
if (!managedLedger.maybeUpdateCursorBeforeTrimmingConsumedLedger()){
Comment thread
Nicklee007 marked this conversation as resolved.
managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE);
}
}
);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,16 +22,21 @@
import static org.testng.Assert.assertNotEquals;
import static org.testng.Assert.assertNotNull;

import java.lang.reflect.Method;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import lombok.Cleanup;
import org.apache.bookkeeper.mledger.ManagedLedgerConfig;
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl;
import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl;
import org.apache.pulsar.broker.service.persistent.PersistentMessageExpiryMonitor;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.Producer;
import org.awaitility.Awaitility;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.testng.Assert;
Expand Down Expand Up @@ -179,4 +184,65 @@ public void testConsumedLedgersTrimNoSubscriptions() throws Exception {
assertEquals(messageIdAfterTrim, MessageId.earliest);

}

@Test
public void testExpiredLedgerDeletionAfterExpiredMessageChecked() throws Exception {
super.baseSetup();
final String ledgerAndCursorName = "testExpiredLedgerDeletionAfterExpiredMessageChecked";
final int totalEntries = 10;
final int ttlSeconds = 1;

final String topicName = "persistent://prop/ns-abc/testExpiredLedgerDeletionAfterExpiredMessageChecked";

@Cleanup
Producer<byte[]> producer = pulsarClient.newProducer()
.topic(topicName)
.producerName("producer-name")
.create();

//set retention parameters, the ledgers are to be deleted as soon as possible
PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get();
ManagedLedgerConfig managedLedgerConfig = persistentTopic.getManagedLedger().getConfig();
managedLedgerConfig.setRetentionSizeInMB(10);
managedLedgerConfig.setRetentionTime(1, TimeUnit.SECONDS);
managedLedgerConfig.setMaxEntriesPerLedger(1000);
managedLedgerConfig.setMinimumRolloverTime(1, TimeUnit.MILLISECONDS);
managedLedgerConfig.setMaximumRolloverTime(50, TimeUnit.MILLISECONDS);

//create managedLedger and managedCursor
ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger();
ManagedCursorImpl managedCursor = (ManagedCursorImpl) managedLedger.openCursor(ledgerAndCursorName);

// write some messages
for (int i = 0; i < totalEntries; i++) {
producer.send(("msg" + i).getBytes());
}

//make sure that all entries should be deleted
Thread.sleep(TimeUnit.SECONDS.toMillis(ttlSeconds));
managedLedger.rollCurrentLedgerIfFull();

PersistentMessageExpiryMonitor monitor = new PersistentMessageExpiryMonitor(topicName, managedCursor.getName(), managedCursor, null);
Position previousMarkDelete = null;
for (int i = 0; i < totalEntries; i++) {
monitor.expireMessages(1);
Position previousPos = previousMarkDelete;
retryStrategically(
(test) -> managedCursor.getMarkDeletedPosition() != null && !managedCursor.getMarkDeletedPosition().equals(previousPos),
5, 100);
previousMarkDelete = managedCursor.getMarkDeletedPosition();
}

//invoke checkConsumedLedgers() to clean the expired ledgers
Method checkConsumedLedgers = BrokerService.class.getDeclaredMethod("checkConsumedLedgers");
checkConsumedLedgers.setAccessible(true);
checkConsumedLedgers.invoke(pulsar.getBrokerService());

Awaitility.await().untilAsserted(()
-> assertEquals(managedLedger.getLedgersInfo().size(), 1));
assertEquals(managedLedger.getLedgersInfo().firstEntry().getValue().getEntries(), 0);

managedCursor.close();
managedLedger.close();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -328,6 +328,11 @@ public void asyncSetProperties(Map<String, String> properties,

}

@Override
public boolean maybeUpdateCursorBeforeTrimmingConsumedLedger() {
return false;
}

@Override
public void trimConsumedLedgersInBackground(CompletableFuture<?> promise) {

Expand Down