From 8965221524f1eda0fb67cf0557d71e7c24650cbe Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 5 Feb 2025 21:33:57 +0800 Subject: [PATCH 01/17] [improve] [broker] Make the estimated entry size more accurate --- .../mledger/impl/ManagedCursorImpl.java | 56 +++++++++++++------ .../impl/cache/RangeEntryCacheImpl.java | 23 ++------ .../InflightReadsLimiterIntegrationTest.java | 15 ++--- 3 files changed, 49 insertions(+), 45 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 203d48933f0a5..5800312db5345 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -3810,26 +3810,50 @@ public int applyMaxSizeCap(int maxEntries, long maxSizeBytes) { if (maxSizeBytes == NO_MAX_SIZE_LIMIT) { return maxEntries; } + int maxEntriesBasedOnSize = + Long.valueOf(estimateEntryCountBySize(maxSizeBytes, readPosition, ledger)).intValue(); + return Math.min(maxEntriesBasedOnSize, maxEntries); + } - double avgEntrySize = ledger.getStats().getEntrySizeAverage(); - if (!Double.isFinite(avgEntrySize)) { - // We don't have yet any stats on the topic entries. Let's try to use the cursor avg size stats - avgEntrySize = (double) entriesReadSize / (double) entriesReadCount; - } - - if (!Double.isFinite(avgEntrySize)) { - // If we still don't have any information, it means this is the first time we attempt reading - // and there are no writes. Let's start with 1 to avoid any overflow and start the avg stats - return 1; + private static long estimateEntryCountBySize(long bytesSize, Position readPosition, ManagedLedgerImpl ml) { + Position posToRead = readPosition; + if (!ml.isValidPosition(readPosition)) { + posToRead = ml.getNextValidPosition(readPosition); } + long result = 0; + long remainingBytesSize = bytesSize; - int maxEntriesBasedOnSize = (int) (maxSizeBytes / avgEntrySize); - if (maxEntriesBasedOnSize < 1) { - // We need to read at least one entry - return 1; + while (remainingBytesSize > 0) { + // Last ledger. + if (posToRead.getLedgerId() == ml.currentLedger.getId()) { + if (ml.currentLedgerSize == 0 || ml.currentLedgerEntries == 0) { + // Only read 1 entry if no entries to read. + return 1; + } + long avg = Math.max(1, ml.currentLedgerSize / ml.currentLedgerEntries); + result += remainingBytesSize / avg; + break; + } + // Skip empty ledger. + LedgerInfo ledgerInfo = ml.getLedgersInfo().get(posToRead.getLedgerId()); + if (ledgerInfo.getSize() == 0 || ledgerInfo.getEntries() == 0) { + posToRead = ml.getNextValidPosition(PositionFactory.create(posToRead.getLedgerId(), Long.MAX_VALUE)); + continue; + } + // Calculate entries by average of ledgers. + long avg = Math.max(1, ledgerInfo.getSize() / ledgerInfo.getEntries()); + long remainEntriesOfLedger = ledgerInfo.getEntries() - posToRead.getEntryId() + 1; + if (remainEntriesOfLedger * avg >= remainingBytesSize) { + result += remainingBytesSize / avg; + break; + } else { + // Calculate for the next ledger. + result += remainEntriesOfLedger / avg; + remainingBytesSize -= remainEntriesOfLedger; + posToRead = ml.getNextValidPosition(PositionFactory.create(posToRead.getLedgerId(), Long.MAX_VALUE)); + } } - - return Math.min(maxEntriesBasedOnSize, maxEntries); + return Math.max(result, 1); // TODO 告诉 RangeEntryCache 这次申请的预估出来的 permits。 } @Override diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index b81015ea63988..dbfa3b348d336 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -71,9 +71,6 @@ public class RangeEntryCacheImpl implements EntryCache { private static final double MB = 1024 * 1024; - private final LongAdder totalAddedEntriesSize = new LongAdder(); - private final LongAdder totalAddedEntriesCount = new LongAdder(); - public RangeEntryCacheImpl(RangeEntryCacheManagerImpl manager, ManagedLedgerImpl ml, boolean copyEntries) { this.manager = manager; this.ml = ml; @@ -152,8 +149,6 @@ public boolean insert(EntryImpl entry) { EntryImpl cacheEntry = EntryImpl.create(position, cachedData); cachedData.release(); if (entries.put(position, cacheEntry)) { - totalAddedEntriesSize.add(entryLength); - totalAddedEntriesCount.increment(); manager.entryAdded(entryLength); return true; } else { @@ -303,7 +298,7 @@ void asyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Position doAsyncReadEntriesByPosition(lh, firstPosition, lastPosition, numberOfEntries, shouldCacheEntry, originalCallback, ctx); } else { - long estimatedEntrySize = getEstimatedEntrySize(); + long estimatedEntrySize = getEstimatedEntrySize(lh); long estimatedReadSize = numberOfEntries * estimatedEntrySize; if (log.isDebugEnabled()) { log.debug("Estimated read size: {} bytes for {} entries with {} estimated entry size", @@ -412,25 +407,15 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio cachedEntries.forEach(entry -> entry.release()); } - // Read all the entries from bookkeeper + // Read all the entries from bookkeeper TODO 优化为只读部分。 pendingReadsManager.readEntries(lh, firstPosition.getEntryId(), lastPosition.getEntryId(), shouldCacheEntry, callback, ctx); } } @VisibleForTesting - public long getEstimatedEntrySize() { - long estimatedEntrySize = getAvgEntrySize(); - if (estimatedEntrySize == 0) { - estimatedEntrySize = DEFAULT_ESTIMATED_ENTRY_SIZE; - } - return estimatedEntrySize + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; - } - - private long getAvgEntrySize() { - long totalAddedEntriesCount = this.totalAddedEntriesCount.sum(); - long totalAddedEntriesSize = this.totalAddedEntriesSize.sum(); - return totalAddedEntriesCount != 0 ? totalAddedEntriesSize / totalAddedEntriesCount : 0; + public long getEstimatedEntrySize(ReadHandle lh) { + return Math.max(1, lh.getLength() / lh.getLastAddConfirmed()); } /** diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/InflightReadsLimiterIntegrationTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/InflightReadsLimiterIntegrationTest.java index 48f0cf08ddff4..6676baf8b555a 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/InflightReadsLimiterIntegrationTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/InflightReadsLimiterIntegrationTest.java @@ -141,10 +141,9 @@ public void testPreciseLimitation(String missingCase) throws Exception { SimpleReadEntriesCallback cb0 = new SimpleReadEntriesCallback(); entryCache.asyncReadEntry(spyCurrentLedger, 125, 125, true, cb0, ctx); cb0.entries.join(); - Long sizePerEntry1 = entryCache.getEstimatedEntrySize(); - Assert.assertEquals(sizePerEntry1, 1 + RangeEntryCacheImpl.BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); + int sizePerEntry = Long.valueOf(entryCache.getEstimatedEntrySize(ml.currentLedger)).intValue(); Awaitility.await().untilAsserted(() -> { - long remainingBytes =limiter.getRemainingBytes(); + long remainingBytes = limiter.getRemainingBytes(); Assert.assertEquals(remainingBytes, totalCapacity); }); log.info("remainingBytes 0: {}", limiter.getRemainingBytes()); @@ -165,7 +164,7 @@ public void testPreciseLimitation(String missingCase) throws Exception { entryCache.asyncReadEntry(spyCurrentLedger, start2, end2, true, cb2, ctx); }).start(); - long bytesAcquired1 = calculateBytesSizeBeforeFirstReading(readCount1 + readCount2, 1); + long bytesAcquired1 = calculateBytesSizeBeforeFirstReading(readCount1 + readCount2, sizePerEntry); long remainingBytesExpected1 = totalCapacity - bytesAcquired1; log.info("acquired : {}", bytesAcquired1); log.info("remainingBytesExpected 0 : {}", remainingBytesExpected1); @@ -178,9 +177,7 @@ public void testPreciseLimitation(String missingCase) throws Exception { Thread.sleep(3000); readCompleteSignal1.countDown(); cb1.entries.join(); - Long sizePerEntry2 = entryCache.getEstimatedEntrySize(); - Assert.assertEquals(sizePerEntry2, 1 + RangeEntryCacheImpl.BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); - long bytesAcquired2 = calculateBytesSizeBeforeFirstReading(readCount2, 1); + long bytesAcquired2 = calculateBytesSizeBeforeFirstReading(readCount2, sizePerEntry); long remainingBytesExpected2 = totalCapacity - bytesAcquired2; log.info("acquired : {}", bytesAcquired2); log.info("remainingBytesExpected 1: {}", remainingBytesExpected2); @@ -191,8 +188,6 @@ public void testPreciseLimitation(String missingCase) throws Exception { readCompleteSignal2.countDown(); cb2.entries.join(); - Long sizePerEntry3 = entryCache.getEstimatedEntrySize(); - Assert.assertEquals(sizePerEntry3, 1 + RangeEntryCacheImpl.BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); Awaitility.await().untilAsserted(() -> { long remainingBytes = limiter.getRemainingBytes(); log.info("remainingBytes 2: {}", remainingBytes); @@ -204,7 +199,7 @@ public void testPreciseLimitation(String missingCase) throws Exception { } private long calculateBytesSizeBeforeFirstReading(int entriesCount, int perEntrySize) { - return entriesCount * (perEntrySize + RangeEntryCacheImpl.BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); + return entriesCount * perEntrySize; } class SimpleReadEntriesCallback implements AsyncCallbacks.ReadEntriesCallback { From dab4ef45b0516022ab6e957db8df986af21c42b1 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 5 Feb 2025 21:56:05 +0800 Subject: [PATCH 02/17] - --- .../bookkeeper/mledger/impl/ManagedCursorImpl.java | 2 +- .../mledger/impl/cache/RangeEntryCacheImpl.java | 13 ++++++------- 2 files changed, 7 insertions(+), 8 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 5800312db5345..1477d6e09981f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -3853,7 +3853,7 @@ private static long estimateEntryCountBySize(long bytesSize, Position readPositi posToRead = ml.getNextValidPosition(PositionFactory.create(posToRead.getLedgerId(), Long.MAX_VALUE)); } } - return Math.max(result, 1); // TODO 告诉 RangeEntryCache 这次申请的预估出来的 permits。 + return Math.max(result, 1); } @Override diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index dbfa3b348d336..8cc35d19cac46 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -55,11 +55,6 @@ */ public class RangeEntryCacheImpl implements EntryCache { - /** - * Overhead per-entry to take into account the envelope. - */ - public static final long BOOKKEEPER_READ_OVERHEAD_PER_ENTRY = 64; - private static final int DEFAULT_ESTIMATED_ENTRY_SIZE = 10 * 1024; private static final boolean DEFAULT_CACHE_INDIVIDUAL_READ_ENTRY = false; private final RangeEntryCacheManagerImpl manager; @@ -407,7 +402,7 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio cachedEntries.forEach(entry -> entry.release()); } - // Read all the entries from bookkeeper TODO 优化为只读部分。 + // Read all the entries from bookkeeper pendingReadsManager.readEntries(lh, firstPosition.getEntryId(), lastPosition.getEntryId(), shouldCacheEntry, callback, ctx); } @@ -415,7 +410,11 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio @VisibleForTesting public long getEstimatedEntrySize(ReadHandle lh) { - return Math.max(1, lh.getLength() / lh.getLastAddConfirmed()); + if (lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { + // No entries stored. + return 1; + } + return Math.max(1, lh.getLength() / (lh.getLastAddConfirmed() + 1)); } /** From 422686a4f4fbb6f8e86101a1e986cd15d709c2ef Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 5 Feb 2025 22:52:49 +0800 Subject: [PATCH 03/17] improve the logic --- .../impl/cache/RangeEntryCacheImpl.java | 20 +++++++++++++++++-- 1 file changed, 18 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index 8cc35d19cac46..3bcc6fd6abc66 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -55,6 +55,11 @@ */ public class RangeEntryCacheImpl implements EntryCache { + /** + * Overhead per-entry to take into account the envelope. + */ + public static final long BOOKKEEPER_READ_OVERHEAD_PER_ENTRY = 64; + private static final int DEFAULT_ESTIMATED_ENTRY_SIZE = 10 * 1024; private static final boolean DEFAULT_CACHE_INDIVIDUAL_READ_ENTRY = false; private final RangeEntryCacheManagerImpl manager; @@ -66,6 +71,9 @@ public class RangeEntryCacheImpl implements EntryCache { private static final double MB = 1024 * 1024; + private final LongAdder totalAddedEntriesSize = new LongAdder(); + private final LongAdder totalAddedEntriesCount = new LongAdder(); + public RangeEntryCacheImpl(RangeEntryCacheManagerImpl manager, ManagedLedgerImpl ml, boolean copyEntries) { this.manager = manager; this.ml = ml; @@ -144,6 +152,8 @@ public boolean insert(EntryImpl entry) { EntryImpl cacheEntry = EntryImpl.create(position, cachedData); cachedData.release(); if (entries.put(position, cacheEntry)) { + totalAddedEntriesSize.add(entryLength); + totalAddedEntriesCount.increment(); manager.entryAdded(entryLength); return true; } else { @@ -410,13 +420,19 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio @VisibleForTesting public long getEstimatedEntrySize(ReadHandle lh) { - if (lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { + if (!lh.isClosed() || lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { // No entries stored. - return 1; + return Math.max(getAvgEntrySize(), DEFAULT_ESTIMATED_ENTRY_SIZE) + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; } return Math.max(1, lh.getLength() / (lh.getLastAddConfirmed() + 1)); } + private long getAvgEntrySize() { + long totalAddedEntriesCount = this.totalAddedEntriesCount.sum(); + long totalAddedEntriesSize = this.totalAddedEntriesSize.sum(); + return totalAddedEntriesCount != 0 ? totalAddedEntriesSize / totalAddedEntriesCount : 0; + } + /** * Reads the entries from Storage. * @param lh the handle From 9bc00c4ee81d8bdf75db8ed95c91ac7e713026d0 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 5 Feb 2025 22:54:30 +0800 Subject: [PATCH 04/17] improve the logic --- .../bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index 3bcc6fd6abc66..9528376ca1624 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -420,7 +420,8 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio @VisibleForTesting public long getEstimatedEntrySize(ReadHandle lh) { - if (!lh.isClosed() || lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { + if ((!lh.isClosed() && lh.getLastAddConfirmed() < 1000) + || lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { // No entries stored. return Math.max(getAvgEntrySize(), DEFAULT_ESTIMATED_ENTRY_SIZE) + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; } From c764c58298306a724755c4d2c8d0b0a4d8617d8f Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 5 Feb 2025 22:57:19 +0800 Subject: [PATCH 05/17] improve the logic --- .../bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index 9528376ca1624..fd1d7977071d6 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -420,8 +420,7 @@ void doAsyncReadEntriesByPosition(ReadHandle lh, Position firstPosition, Positio @VisibleForTesting public long getEstimatedEntrySize(ReadHandle lh) { - if ((!lh.isClosed() && lh.getLastAddConfirmed() < 1000) - || lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { + if (lh.getLength() == 0 || lh.getLastAddConfirmed() < 0) { // No entries stored. return Math.max(getAvgEntrySize(), DEFAULT_ESTIMATED_ENTRY_SIZE) + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; } From 23ac2e127b8e8def3aa8f17e7a96ea640db55816 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 11:14:13 +0800 Subject: [PATCH 06/17] add test --- .../mledger/impl/ManagedCursorTest.java | 74 +++++++++++++++++++ 1 file changed, 74 insertions(+) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index d3ea98131ad8f..f4c381dad0b6d 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -5164,6 +5164,80 @@ public void findEntryFailed(ManagedLedgerException exception, Optional assertEquals(positionRef4.get(), position4); } + @Test + public void testEstimateEntryCountBySize() throws Exception { + final String mlName = "ml-" + UUID.randomUUID().toString().replaceAll("-", ""); + ManagedLedgerImpl ml = (ManagedLedgerImpl) factory.open(mlName); + // Verify: no entry to read + long entryCount0 = + ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ml.currentLedger.getId(), 0), ml); + assertEquals(entryCount0, 1); + // Avoid trimming ledgers. + ml.openCursor("c1"); + + // Build data. + for (int i = 0; i < 100; i++) { + ml.addEntry(new byte[]{1}); + } + long ledger1 = ml.currentLedger.getId(); + ml.currentLedger.close(); + ml.ledgerClosed(ml.currentLedger); + for (int i = 0; i < 100; i++) { + ml.addEntry(new byte[]{1, 2}); + } + long ledger2 = ml.currentLedger.getId(); + ml.currentLedger.close(); + ml.ledgerClosed(ml.currentLedger); + for (int i = 0; i < 100; i++) { + ml.addEntry(new byte[]{1, 2, 3, 4}); + } + long ledger3 = ml.currentLedger.getId(); + MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo1 = ml.getLedgersInfo().get(ledger1); + MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo2 = ml.getLedgersInfo().get(ledger2); + long average1 = ledgerInfo1.getSize() / ledgerInfo1.getEntries(); + long average2 = ledgerInfo2.getSize() / ledgerInfo2.getEntries(); + long average3 = ml.currentLedgerSize / ml.currentLedgerEntries; + assertEquals(average1, 1); + assertEquals(average2, 2); + assertEquals(average3, 4); + + // Test: the individual ledgers. + long entryCount1 = + ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ledger1, 0), ml); + assertEquals(entryCount1, 16); + long entryCount2 = + ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ledger2, 0), ml); + assertEquals(entryCount2, 8); + long entryCount3 = + ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ledger3, 0), ml); + assertEquals(entryCount3, 4); + + // Test: across ledgers. + long entryCount4 = + ManagedCursorImpl.estimateEntryCountBySize(116, PositionFactory.create(ledger1, 0), ml); + assertEquals(entryCount4, 108); + long entryCount5 = + ManagedCursorImpl.estimateEntryCountBySize(216, PositionFactory.create(ledger2, 0), ml); + assertEquals(entryCount5, 104); + long entryCount6 = + ManagedCursorImpl.estimateEntryCountBySize(316, PositionFactory.create(ledger1, 0), ml); + assertEquals(entryCount6, 204); + + long entryCount7 = + ManagedCursorImpl.estimateEntryCountBySize(36, PositionFactory.create(ledger1, 80), ml); + assertEquals(entryCount7, 28); + long entryCount8 = + ManagedCursorImpl.estimateEntryCountBySize(56, PositionFactory.create(ledger2, 80), ml); + assertEquals(entryCount8, 24); + long entryCount9 = + ManagedCursorImpl.estimateEntryCountBySize(236, PositionFactory.create(ledger1, 80), ml); + assertEquals(entryCount9, 124); + + + // cleanup. + ml.delete(); + } + @Test void testForceCursorRecovery() throws Exception { TestPulsarMockBookKeeper bk = new TestPulsarMockBookKeeper(executor); From 72896007b3e466753d52fada0b03409f6e51e110 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 11:49:06 +0800 Subject: [PATCH 07/17] add test --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 1 + 1 file changed, 1 insertion(+) 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 7426059e576f6..7996ca4b63b97 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 @@ -224,6 +224,7 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { private final CallbackMutex offloadMutex = new CallbackMutex(); public static final CompletableFuture NULL_OFFLOAD_PROMISE = CompletableFuture .completedFuture(PositionFactory.LATEST); + @VisibleForTesting protected volatile LedgerHandle currentLedger; protected volatile long currentLedgerEntries = 0; protected volatile long currentLedgerSize = 0; From c27e0c763834fb3a41ca9f95d114d97035c9632a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 12:01:49 +0800 Subject: [PATCH 08/17] add test --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 1 + 1 file changed, 1 insertion(+) 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 7996ca4b63b97..c890ba01f634e 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 @@ -225,6 +225,7 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { public static final CompletableFuture NULL_OFFLOAD_PROMISE = CompletableFuture .completedFuture(PositionFactory.LATEST); @VisibleForTesting + @Getter protected volatile LedgerHandle currentLedger; protected volatile long currentLedgerEntries = 0; protected volatile long currentLedgerSize = 0; From 9ee3410e45869a8b8ac8d3461a7b63337c74849e Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 12:31:32 +0800 Subject: [PATCH 09/17] add test --- .../mledger/impl/ManagedCursorImpl.java | 14 +++++++------- .../mledger/impl/ManagedCursorTest.java | 18 +++++++++--------- 2 files changed, 16 insertions(+), 16 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 1477d6e09981f..86aab6152f6b8 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -3815,7 +3815,7 @@ public int applyMaxSizeCap(int maxEntries, long maxSizeBytes) { return Math.min(maxEntriesBasedOnSize, maxEntries); } - private static long estimateEntryCountBySize(long bytesSize, Position readPosition, ManagedLedgerImpl ml) { + static long estimateEntryCountBySize(long bytesSize, Position readPosition, ManagedLedgerImpl ml) { Position posToRead = readPosition; if (!ml.isValidPosition(readPosition)) { posToRead = ml.getNextValidPosition(readPosition); @@ -3825,12 +3825,12 @@ private static long estimateEntryCountBySize(long bytesSize, Position readPositi while (remainingBytesSize > 0) { // Last ledger. - if (posToRead.getLedgerId() == ml.currentLedger.getId()) { - if (ml.currentLedgerSize == 0 || ml.currentLedgerEntries == 0) { + if (posToRead.getLedgerId() == ml.getCurrentLedger().getId()) { + if (ml.getCurrentLedgerSize() == 0 || ml.getCurrentLedgerEntries() == 0) { // Only read 1 entry if no entries to read. return 1; } - long avg = Math.max(1, ml.currentLedgerSize / ml.currentLedgerEntries); + long avg = Math.max(1, ml.getCurrentLedgerSize() / ml.getCurrentLedgerEntries()); result += remainingBytesSize / avg; break; } @@ -3842,14 +3842,14 @@ private static long estimateEntryCountBySize(long bytesSize, Position readPositi } // Calculate entries by average of ledgers. long avg = Math.max(1, ledgerInfo.getSize() / ledgerInfo.getEntries()); - long remainEntriesOfLedger = ledgerInfo.getEntries() - posToRead.getEntryId() + 1; + long remainEntriesOfLedger = ledgerInfo.getEntries() - posToRead.getEntryId(); if (remainEntriesOfLedger * avg >= remainingBytesSize) { result += remainingBytesSize / avg; break; } else { // Calculate for the next ledger. - result += remainEntriesOfLedger / avg; - remainingBytesSize -= remainEntriesOfLedger; + result += remainEntriesOfLedger; + remainingBytesSize -= remainEntriesOfLedger * avg; posToRead = ml.getNextValidPosition(PositionFactory.create(posToRead.getLedgerId(), Long.MAX_VALUE)); } } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index f4c381dad0b6d..023f3d96da390 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -5169,8 +5169,9 @@ public void testEstimateEntryCountBySize() throws Exception { final String mlName = "ml-" + UUID.randomUUID().toString().replaceAll("-", ""); ManagedLedgerImpl ml = (ManagedLedgerImpl) factory.open(mlName); // Verify: no entry to read + long entryCount0 = - ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ml.currentLedger.getId(), 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ml.getCurrentLedger().getId(), 0), ml); assertEquals(entryCount0, 1); // Avoid trimming ledgers. ml.openCursor("c1"); @@ -5179,19 +5180,19 @@ public void testEstimateEntryCountBySize() throws Exception { for (int i = 0; i < 100; i++) { ml.addEntry(new byte[]{1}); } - long ledger1 = ml.currentLedger.getId(); - ml.currentLedger.close(); - ml.ledgerClosed(ml.currentLedger); + long ledger1 = ml.getCurrentLedger().getId(); + ml.getCurrentLedger().close(); + ml.ledgerClosed(ml.getCurrentLedger()); for (int i = 0; i < 100; i++) { ml.addEntry(new byte[]{1, 2}); } - long ledger2 = ml.currentLedger.getId(); - ml.currentLedger.close(); - ml.ledgerClosed(ml.currentLedger); + long ledger2 = ml.getCurrentLedger().getId(); + ml.getCurrentLedger().close(); + ml.ledgerClosed(ml.getCurrentLedger()); for (int i = 0; i < 100; i++) { ml.addEntry(new byte[]{1, 2, 3, 4}); } - long ledger3 = ml.currentLedger.getId(); + long ledger3 = ml.getCurrentLedger().getId(); MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo1 = ml.getLedgersInfo().get(ledger1); MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo2 = ml.getLedgersInfo().get(ledger2); long average1 = ledgerInfo1.getSize() / ledgerInfo1.getEntries(); @@ -5233,7 +5234,6 @@ public void testEstimateEntryCountBySize() throws Exception { ManagedCursorImpl.estimateEntryCountBySize(236, PositionFactory.create(ledger1, 80), ml); assertEquals(entryCount9, 124); - // cleanup. ml.delete(); } From daebb6bfadd50bfb76aa1a9b6d6ea4ce497fd70a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 12:32:22 +0800 Subject: [PATCH 10/17] add test --- .../org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 023f3d96da390..51b577c678667 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -5197,7 +5197,7 @@ public void testEstimateEntryCountBySize() throws Exception { MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo2 = ml.getLedgersInfo().get(ledger2); long average1 = ledgerInfo1.getSize() / ledgerInfo1.getEntries(); long average2 = ledgerInfo2.getSize() / ledgerInfo2.getEntries(); - long average3 = ml.currentLedgerSize / ml.currentLedgerEntries; + long average3 = ml.getCurrentLedgerSize() / ml.getCurrentLedgerEntries(); assertEquals(average1, 1); assertEquals(average2, 2); assertEquals(average3, 4); From 1b5f10d75f5bfdf453c43552a643f863ce8c724b Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 14:45:56 +0800 Subject: [PATCH 11/17] fix tests --- .../apache/bookkeeper/mledger/impl/ManagedCursorTest.java | 5 +++-- .../org/apache/pulsar/broker/stats/ConsumerStatsTest.java | 4 ++++ 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 51b577c678667..ec7d61d903bd5 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -689,8 +689,9 @@ void testAsyncReadWithMaxSizeByte() throws Exception { ledger.addEntry(new byte[1024]); } - // First time, since we don't have info, we'll get 1 single entry - readAndCheck(cursor, 10, 3 * 1024, 1); + // Since https://github.com/apache/pulsar/pull/23931 improved the performance of delivery, the consumer + // will get more messages than before(it only receives 1 messages at the first delivery), + readAndCheck(cursor, 10, 3 * 1024, 3); // We should only return 3 entries, based on the max size readAndCheck(cursor, 20, 3 * 1024, 3); // If maxSize is < avg, we should get 1 entry diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java index 59a911500e5d9..2e1bb159d51c4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java @@ -480,6 +480,10 @@ public void testAvgMessagesPerEntry() throws Exception { metadataConsumer.put("matchValueReschedule", "producer2"); @Cleanup Consumer consumer = pulsarClient.newConsumer(Schema.STRING).topic(topic).properties(metadataConsumer) + // Since https://github.com/apache/pulsar/pull/23931 improved the performance of delivery, the consumer + // will get more messages than before(it only receives 1 messages at the first delivery), we set queue + // size to `1` to keep the test passing. + .receiverQueueSize(1) .subscriptionName(subName).subscriptionInitialPosition(SubscriptionInitialPosition.Earliest).subscribe(); int counter = 0; From 30646c0dfecc178de3d805b4e242dd8623bd4cde Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 6 Feb 2025 15:14:54 +0800 Subject: [PATCH 12/17] fix test --- .../mledger/impl/ManagedCursorTest.java | 22 ++++++++++++------- 1 file changed, 14 insertions(+), 8 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index ec7d61d903bd5..ec625b7dfdfc2 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -3915,9 +3915,10 @@ public void testReadEntriesOrWaitWithMaxSize() throws Exception { ledger.addEntry(new byte[1024]); } - // First time, since we don't have info, we'll get 1 single entry + // Since https://github.com/apache/pulsar/pull/23931 improved the performance of delivery, the consumer + // will get more messages than before(it only receives 1 messages at the first delivery), List entries = c.readEntriesOrWait(10, 3 * 1024); - assertEquals(entries.size(), 1); + assertEquals(entries.size(), 3); entries.forEach(Entry::release); // We should only return 3 entries, based on the max size @@ -5216,25 +5217,30 @@ public void testEstimateEntryCountBySize() throws Exception { // Test: across ledgers. long entryCount4 = - ManagedCursorImpl.estimateEntryCountBySize(116, PositionFactory.create(ledger1, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(100 + 16, PositionFactory.create(ledger1, 0), ml); assertEquals(entryCount4, 108); long entryCount5 = - ManagedCursorImpl.estimateEntryCountBySize(216, PositionFactory.create(ledger2, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(200 + 16, PositionFactory.create(ledger2, 0), ml); assertEquals(entryCount5, 104); long entryCount6 = - ManagedCursorImpl.estimateEntryCountBySize(316, PositionFactory.create(ledger1, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(100 + 200 + 16, PositionFactory.create(ledger1, 0), ml); assertEquals(entryCount6, 204); long entryCount7 = - ManagedCursorImpl.estimateEntryCountBySize(36, PositionFactory.create(ledger1, 80), ml); + ManagedCursorImpl.estimateEntryCountBySize(20 + 16, PositionFactory.create(ledger1, 80), ml); assertEquals(entryCount7, 28); long entryCount8 = - ManagedCursorImpl.estimateEntryCountBySize(56, PositionFactory.create(ledger2, 80), ml); + ManagedCursorImpl.estimateEntryCountBySize(40 + 16, PositionFactory.create(ledger2, 80), ml); assertEquals(entryCount8, 24); long entryCount9 = - ManagedCursorImpl.estimateEntryCountBySize(236, PositionFactory.create(ledger1, 80), ml); + ManagedCursorImpl.estimateEntryCountBySize(20 + 200 + 16, PositionFactory.create(ledger1, 80), ml); assertEquals(entryCount9, 124); + // Test: read more than entries written. + long entryCount10 = + ManagedCursorImpl.estimateEntryCountBySize(100 + 200 + 400 + 16, PositionFactory.create(ledger1, 0), ml); + assertEquals(entryCount10, 304); + // cleanup. ml.delete(); } From 833ea8ebf4565902b6eaa39b20e3ed59b5ab90a6 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 11 Feb 2025 15:12:02 +0800 Subject: [PATCH 13/17] fix test --- .../service/BatchMessageWithBatchIndexLevelTest.java | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageWithBatchIndexLevelTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageWithBatchIndexLevelTest.java index 7fa7bf078e0c5..f21ac130e3cfd 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageWithBatchIndexLevelTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BatchMessageWithBatchIndexLevelTest.java @@ -85,7 +85,7 @@ public void testBatchMessageAck() { .newConsumer() .topic(topicName) .subscriptionName(subscriptionName) - .receiverQueueSize(10) + .receiverQueueSize(50) .subscriptionType(SubscriptionType.Shared) .enableBatchIndexAcknowledgment(true) .negativeAckRedeliveryDelay(100, TimeUnit.MILLISECONDS) @@ -114,27 +114,29 @@ public void testBatchMessageAck() { consumer.acknowledge(receive1); consumer.acknowledge(receive2); Awaitility.await().untilAsserted(() -> { - assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 18); + // Since https://github.com/apache/pulsar/pull/23931 improved the mechanism of estimate average entry size, + // broker will deliver much messages than before. So edit 18 -> 38 here. + assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 38); }); Message receive3 = consumer.receive(); Message receive4 = consumer.receive(); consumer.acknowledge(receive3); consumer.acknowledge(receive4); Awaitility.await().untilAsserted(() -> { - assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 16); + assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 36); }); // Block cmd-flow send until verify finish. see: https://github.com/apache/pulsar/pull/17436. consumer.pause(); Message receive5 = consumer.receive(); consumer.negativeAcknowledge(receive5); Awaitility.await().pollInterval(1, TimeUnit.MILLISECONDS).untilAsserted(() -> { - assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 0); + assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 20); }); // Unblock cmd-flow. consumer.resume(); consumer.receive(); Awaitility.await().untilAsserted(() -> { - assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 16); + assertEquals(dispatcher.getConsumers().get(0).getUnackedMessages(), 36); }); } From d4cd9fb05d8be8a9d5d0f1f48e1849814b65e0fb Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 21 Feb 2025 16:06:23 +0800 Subject: [PATCH 14/17] add BOOKKEEPER_READ_OVERHEAD_PER_ENTRY --- .../mledger/impl/ManagedCursorImpl.java | 6 ++-- .../impl/cache/RangeEntryCacheImpl.java | 2 +- .../mledger/impl/ManagedCursorTest.java | 35 +++++++++---------- 3 files changed, 22 insertions(+), 21 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 86aab6152f6b8..32d46ff1c3c40 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -24,6 +24,7 @@ import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.DEFAULT_LEDGER_DELETE_BACKOFF_TIME_SEC; import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.DEFAULT_LEDGER_DELETE_RETRIES; import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.createManagedLedgerException; +import static org.apache.bookkeeper.mledger.impl.cache.RangeEntryCacheImpl.BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; import static org.apache.bookkeeper.mledger.util.Errors.isNoSuchLedgerExistsException; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.MoreObjects; @@ -3830,7 +3831,8 @@ static long estimateEntryCountBySize(long bytesSize, Position readPosition, Mana // Only read 1 entry if no entries to read. return 1; } - long avg = Math.max(1, ml.getCurrentLedgerSize() / ml.getCurrentLedgerEntries()); + long avg = Math.max(1, ml.getCurrentLedgerSize() / ml.getCurrentLedgerEntries()) + + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; result += remainingBytesSize / avg; break; } @@ -3841,7 +3843,7 @@ static long estimateEntryCountBySize(long bytesSize, Position readPosition, Mana continue; } // Calculate entries by average of ledgers. - long avg = Math.max(1, ledgerInfo.getSize() / ledgerInfo.getEntries()); + long avg = Math.max(1, ledgerInfo.getSize() / ledgerInfo.getEntries()) + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; long remainEntriesOfLedger = ledgerInfo.getEntries() - posToRead.getEntryId(); if (remainEntriesOfLedger * avg >= remainingBytesSize) { result += remainingBytesSize / avg; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index fd1d7977071d6..9a2de9ba8c41d 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -424,7 +424,7 @@ public long getEstimatedEntrySize(ReadHandle lh) { // No entries stored. return Math.max(getAvgEntrySize(), DEFAULT_ESTIMATED_ENTRY_SIZE) + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; } - return Math.max(1, lh.getLength() / (lh.getLastAddConfirmed() + 1)); + return Math.max(1, lh.getLength() / (lh.getLastAddConfirmed() + 1)) + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; } private long getAvgEntrySize() { diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index ec625b7dfdfc2..0242e7d19f1c9 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -18,6 +18,7 @@ */ package org.apache.bookkeeper.mledger.impl; +import static org.apache.bookkeeper.mledger.impl.cache.RangeEntryCacheImpl.BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.Mockito.any; import static org.mockito.Mockito.doAnswer; @@ -5170,8 +5171,6 @@ public void findEntryFailed(ManagedLedgerException exception, Optional public void testEstimateEntryCountBySize() throws Exception { final String mlName = "ml-" + UUID.randomUUID().toString().replaceAll("-", ""); ManagedLedgerImpl ml = (ManagedLedgerImpl) factory.open(mlName); - // Verify: no entry to read - long entryCount0 = ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ml.getCurrentLedger().getId(), 0), ml); assertEquals(entryCount0, 1); @@ -5197,48 +5196,48 @@ public void testEstimateEntryCountBySize() throws Exception { long ledger3 = ml.getCurrentLedger().getId(); MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo1 = ml.getLedgersInfo().get(ledger1); MLDataFormats.ManagedLedgerInfo.LedgerInfo ledgerInfo2 = ml.getLedgersInfo().get(ledger2); - long average1 = ledgerInfo1.getSize() / ledgerInfo1.getEntries(); - long average2 = ledgerInfo2.getSize() / ledgerInfo2.getEntries(); - long average3 = ml.getCurrentLedgerSize() / ml.getCurrentLedgerEntries(); - assertEquals(average1, 1); - assertEquals(average2, 2); - assertEquals(average3, 4); + long average1 = ledgerInfo1.getSize() / ledgerInfo1.getEntries() + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; + long average2 = ledgerInfo2.getSize() / ledgerInfo2.getEntries() + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; + long average3 = ml.getCurrentLedgerSize() / ml.getCurrentLedgerEntries() + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY; + assertEquals(average1, 1 + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); + assertEquals(average2, 2 + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); + assertEquals(average3, 4 + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); // Test: the individual ledgers. long entryCount1 = - ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ledger1, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(average1 * 16, PositionFactory.create(ledger1, 0), ml); assertEquals(entryCount1, 16); long entryCount2 = - ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ledger2, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(average2 * 8, PositionFactory.create(ledger2, 0), ml); assertEquals(entryCount2, 8); long entryCount3 = - ManagedCursorImpl.estimateEntryCountBySize(16, PositionFactory.create(ledger3, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize(average3 * 4, PositionFactory.create(ledger3, 0), ml); assertEquals(entryCount3, 4); // Test: across ledgers. long entryCount4 = - ManagedCursorImpl.estimateEntryCountBySize(100 + 16, PositionFactory.create(ledger1, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize((average1 * 100) + (average2 * 8), PositionFactory.create(ledger1, 0), ml); assertEquals(entryCount4, 108); long entryCount5 = - ManagedCursorImpl.estimateEntryCountBySize(200 + 16, PositionFactory.create(ledger2, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize((average2 * 100) + (average3 * 4), PositionFactory.create(ledger2, 0), ml); assertEquals(entryCount5, 104); long entryCount6 = - ManagedCursorImpl.estimateEntryCountBySize(100 + 200 + 16, PositionFactory.create(ledger1, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize((average1 * 100) + (average2 * 100) + (average3 * 4), PositionFactory.create(ledger1, 0), ml); assertEquals(entryCount6, 204); long entryCount7 = - ManagedCursorImpl.estimateEntryCountBySize(20 + 16, PositionFactory.create(ledger1, 80), ml); + ManagedCursorImpl.estimateEntryCountBySize((average1 * 20) + (average2 * 8), PositionFactory.create(ledger1, 80), ml); assertEquals(entryCount7, 28); long entryCount8 = - ManagedCursorImpl.estimateEntryCountBySize(40 + 16, PositionFactory.create(ledger2, 80), ml); + ManagedCursorImpl.estimateEntryCountBySize((average2 * 20) + (average3 * 4), PositionFactory.create(ledger2, 80), ml); assertEquals(entryCount8, 24); long entryCount9 = - ManagedCursorImpl.estimateEntryCountBySize(20 + 200 + 16, PositionFactory.create(ledger1, 80), ml); + ManagedCursorImpl.estimateEntryCountBySize((average1 * 20) + (average2 * 100) + (average3 * 4), PositionFactory.create(ledger1, 80), ml); assertEquals(entryCount9, 124); // Test: read more than entries written. long entryCount10 = - ManagedCursorImpl.estimateEntryCountBySize(100 + 200 + 400 + 16, PositionFactory.create(ledger1, 0), ml); + ManagedCursorImpl.estimateEntryCountBySize((average1 * 100) + (average2 * 100) + (average3 * 100) + (average3 * 4) , PositionFactory.create(ledger1, 0), ml); assertEquals(entryCount10, 304); // cleanup. From 10664c0fc04720d2d6bf45939e4a87b04978bbd1 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 21 Feb 2025 17:17:27 +0800 Subject: [PATCH 15/17] fix test --- .../apache/bookkeeper/mledger/impl/ManagedCursorTest.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 0242e7d19f1c9..8611a247f4a03 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -687,14 +687,15 @@ void testAsyncReadWithMaxSizeByte() throws Exception { ManagedCursor cursor = ledger.openCursor("c1"); for (int i = 0; i < 100; i++) { - ledger.addEntry(new byte[1024]); + ledger.addEntry(new byte[(int) (1024)]); } // Since https://github.com/apache/pulsar/pull/23931 improved the performance of delivery, the consumer // will get more messages than before(it only receives 1 messages at the first delivery), - readAndCheck(cursor, 10, 3 * 1024, 3); + int avg = (int) (BOOKKEEPER_READ_OVERHEAD_PER_ENTRY + 1024); + readAndCheck(cursor, 10, 3 * avg, 3); // We should only return 3 entries, based on the max size - readAndCheck(cursor, 20, 3 * 1024, 3); + readAndCheck(cursor, 20, 3 * avg, 3); // If maxSize is < avg, we should get 1 entry readAndCheck(cursor, 10, 500, 1); } From 39e4e02c5c95fdbf5ce8df9aadc80adbed2a0503 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 21 Feb 2025 19:55:10 +0800 Subject: [PATCH 16/17] fix test --- .../apache/bookkeeper/mledger/impl/ManagedCursorTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 8611a247f4a03..1cb09d995393c 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -3919,12 +3919,13 @@ public void testReadEntriesOrWaitWithMaxSize() throws Exception { // Since https://github.com/apache/pulsar/pull/23931 improved the performance of delivery, the consumer // will get more messages than before(it only receives 1 messages at the first delivery), - List entries = c.readEntriesOrWait(10, 3 * 1024); + int avg = (int) (1024 + BOOKKEEPER_READ_OVERHEAD_PER_ENTRY); + List entries = c.readEntriesOrWait(10, 3 * avg); assertEquals(entries.size(), 3); entries.forEach(Entry::release); // We should only return 3 entries, based on the max size - entries = c.readEntriesOrWait(10, 3 * 1024); + entries = c.readEntriesOrWait(10, 3 * avg); assertEquals(entries.size(), 3); entries.forEach(Entry::release); From a46123fb257d47432f0fe1855dd2f67b6c205a47 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 25 Feb 2025 11:31:05 +0800 Subject: [PATCH 17/17] address comments --- .../broker/stats/ConsumerStatsTest.java | 23 +++++++++++-------- 1 file changed, 14 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java index 2e1bb159d51c4..fc650127f90a8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/stats/ConsumerStatsTest.java @@ -445,8 +445,13 @@ public void testAvgMessagesPerEntry() throws Exception { .batchingMaxPublishDelay(5, TimeUnit.SECONDS) .batchingMaxBytes(Integer.MAX_VALUE) .create(); - - producer.send("first-message"); + // The first messages deliver: 20 msgs. + // Average of "messages per batch" is "1". + for (int i = 0; i < 20; i++) { + producer.send("first-message"); + } + // The second messages deliver: 20 msgs. + // Average of "messages per batch" is "Math.round(1 * 0.9 + 20 * 0.1) = 2.9 ~ 3". List> futures = new ArrayList<>(); for (int i = 0; i < 20; i++) { futures.add(producer.sendAsync("message")); @@ -480,10 +485,7 @@ public void testAvgMessagesPerEntry() throws Exception { metadataConsumer.put("matchValueReschedule", "producer2"); @Cleanup Consumer consumer = pulsarClient.newConsumer(Schema.STRING).topic(topic).properties(metadataConsumer) - // Since https://github.com/apache/pulsar/pull/23931 improved the performance of delivery, the consumer - // will get more messages than before(it only receives 1 messages at the first delivery), we set queue - // size to `1` to keep the test passing. - .receiverQueueSize(1) + .receiverQueueSize(20) .subscriptionName(subName).subscriptionInitialPosition(SubscriptionInitialPosition.Earliest).subscribe(); int counter = 0; @@ -498,14 +500,17 @@ public void testAvgMessagesPerEntry() throws Exception { } } - assertEquals(21, counter); + assertEquals(40, counter); ConsumerStats consumerStats = admin.topics().getStats(topic).getSubscriptions().get(subName).getConsumers().get(0); - assertEquals(21, consumerStats.getMsgOutCounter()); + assertEquals(40, consumerStats.getMsgOutCounter()); - // Math.round(1 * 0.9 + 0.1 * (20 / 1)) + // The first messages deliver: 20 msgs. + // Average of "messages per batch" is "1". + // The second messages deliver: 20 msgs. + // Average of "messages per batch" is "Math.round(1 * 0.9 + 20 * 0.1) = 2.9 ~ 3". int avgMessagesPerEntry = consumerStats.getAvgMessagesPerEntry(); assertEquals(3, avgMessagesPerEntry); }