Skip to content

Commit d76d4cb

Browse files
poorbarcodelhotari
authored andcommitted
[fix][ml] There are two same-named managed ledgers in the one broker (#18688)
(cherry picked from commit d7186a6)
1 parent 1ca7376 commit d76d4cb

3 files changed

Lines changed: 65 additions & 3 deletions

File tree

managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -485,9 +485,19 @@ public void asyncOpenReadOnlyCursor(String managedLedgerName, Position startPosi
485485
}
486486

487487
void close(ManagedLedger ledger) {
488-
// Remove the ledger from the internal factory cache
489-
ledgers.remove(ledger.getName());
490-
entryCacheManager.removeEntryCache(ledger.getName());
488+
// If the future in map is not done or has exceptionally complete, it means that @param-ledger is not in the
489+
// map.
490+
CompletableFuture<ManagedLedgerImpl> ledgerFuture = ledgers.get(ledger.getName());
491+
if (ledgerFuture == null || !ledgerFuture.isDone() || ledgerFuture.isCompletedExceptionally()){
492+
return;
493+
}
494+
if (ledgerFuture.join() != ledger){
495+
return;
496+
}
497+
// Remove the ledger from the internal factory cache.
498+
if (ledgers.remove(ledger.getName(), ledgerFuture)) {
499+
entryCacheManager.removeEntryCache(ledger.getName());
500+
}
491501
}
492502

493503
public CompletableFuture<Void> shutdownAsync() throws ManagedLedgerException {

managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1482,6 +1482,15 @@ private void closeAllCursors(CloseCallback callback, final Object ctx) {
14821482

14831483
@Override
14841484
public synchronized void createComplete(int rc, final LedgerHandle lh, Object ctx) {
1485+
if (STATE_UPDATER.get(this) == State.Closed) {
1486+
if (lh != null) {
1487+
log.warn("[{}] ledger create completed after the managed ledger is closed rc={} ledger={}, so just"
1488+
+ " close this ledger handle.", name, rc, lh != null ? lh.getId() : -1);
1489+
lh.closeAsync();
1490+
}
1491+
return;
1492+
}
1493+
14851494
if (log.isDebugEnabled()) {
14861495
log.debug("[{}] createComplete rc={} ledger={}", name, rc, lh != null ? lh.getId() : -1);
14871496
}

managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryTest.java

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,12 +20,16 @@
2020

2121
import static org.testng.Assert.assertEquals;
2222

23+
import java.util.UUID;
24+
import java.util.concurrent.TimeUnit;
2325
import org.apache.bookkeeper.mledger.ManagedCursor;
2426
import org.apache.bookkeeper.mledger.ManagedLedgerConfig;
2527
import org.apache.bookkeeper.mledger.ManagedLedgerInfo;
2628
import org.apache.bookkeeper.mledger.ManagedLedgerInfo.CursorInfo;
2729
import org.apache.bookkeeper.mledger.ManagedLedgerInfo.MessageRangeInfo;
2830
import org.apache.bookkeeper.test.MockedBookKeeperTestCase;
31+
import org.awaitility.Awaitility;
32+
import org.testng.Assert;
2933
import org.testng.annotations.Test;
3034

3135
public class ManagedLedgerFactoryTest extends MockedBookKeeperTestCase {
@@ -71,4 +75,43 @@ public void testGetManagedLedgerInfoWithClose() throws Exception {
7175
assertEquals(mri.to.entryId, 0);
7276
}
7377

78+
/**
79+
* see: https://github.com/apache/pulsar/pull/18688
80+
*/
81+
@Test
82+
public void testConcurrentCloseLedgerAndSwitchLedgerForReproduceIssue() throws Exception {
83+
String managedLedgerName = "lg_" + UUID.randomUUID().toString().replaceAll("-", "_");
84+
85+
ManagedLedgerConfig config = new ManagedLedgerConfig();
86+
config.setThrottleMarkDelete(1);
87+
config.setMaximumRolloverTime(Integer.MAX_VALUE, TimeUnit.SECONDS);
88+
config.setMaxEntriesPerLedger(5);
89+
90+
// create managedLedger once and close it.
91+
ManagedLedgerImpl managedLedger1 = (ManagedLedgerImpl) factory.open(managedLedgerName, config);
92+
waitManagedLedgerStateEquals(managedLedger1, ManagedLedgerImpl.State.LedgerOpened);
93+
managedLedger1.close();
94+
95+
// create managedLedger the second time.
96+
ManagedLedgerImpl managedLedger2 = (ManagedLedgerImpl) factory.open(managedLedgerName, config);
97+
waitManagedLedgerStateEquals(managedLedger2, ManagedLedgerImpl.State.LedgerOpened);
98+
99+
// Mock the task create ledger complete now, it will change the state to another value which not is Closed.
100+
// Close managedLedger1 the second time.
101+
managedLedger1.createComplete(1, null, null);
102+
managedLedger1.close();
103+
104+
// Verify managedLedger2 is still there.
105+
Assert.assertFalse(factory.ledgers.isEmpty());
106+
Assert.assertEquals(factory.ledgers.get(managedLedger2.getName()).join(), managedLedger2);
107+
108+
// cleanup.
109+
managedLedger2.close();
110+
}
111+
112+
private void waitManagedLedgerStateEquals(ManagedLedgerImpl managedLedger, ManagedLedgerImpl.State expectedStat){
113+
Awaitility.await().untilAsserted(() ->
114+
Assert.assertTrue(managedLedger.getState() == expectedStat));
115+
}
116+
74117
}

0 commit comments

Comments
 (0)