sijie closed pull request #1901: Estimate the size of a managed ledger suffix 
from a position
URL: https://github.com/apache/incubator-pulsar/pull/1901
 
 
   

This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:

As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):

diff --git 
a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java 
b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java
index 186a450c87..584a376fac 100644
--- 
a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java
+++ 
b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java
@@ -528,6 +528,13 @@ void asyncFindNewestMatching(FindPositionConstraint 
constraint, Predicate<Entry>
      */
     int getTotalNonContiguousDeletedMessagesRange();
 
+    /**
+     * Returns the estimated size of the unacknowledged backlog for this cursor
+     *
+     * @return the estimated size from the mark delete position of the cursor
+     */
+    long getEstimatedSizeSinceMarkDeletePosition();
+
     /**
      * Returns cursor throttle mark-delete rate.
      *
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 194e8c0e6f..898ac8b803 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
@@ -660,6 +660,11 @@ public int getTotalNonContiguousDeletedMessagesRange() {
         return individualDeletedMessages.asRanges().size();
     }
 
+    @Override
+    public long getEstimatedSizeSinceMarkDeletePosition() {
+        return ledger.estimateBacklogFromPosition(markDeletePosition);
+    }
+
     @Override
     public long getNumberOfEntriesInBacklog() {
         if (log.isDebugEnabled()) {
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 1a606b7d4a..68bc801a53 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
@@ -930,6 +930,24 @@ public long getEstimatedBacklogSize() {
         }
     }
 
+    long estimateBacklogFromPosition(PositionImpl pos) {
+        synchronized (this) {
+            LedgerInfo ledgerInfo = ledgers.get(pos.getLedgerId());
+            if (ledgerInfo == null) {
+                return getTotalSize(); // position no longer in managed 
ledger, so return total size
+            }
+            long sizeBeforePosLedger = ledgers.values().stream().filter(li -> 
li.getLedgerId() < pos.getLedgerId())
+                .mapToLong(li -> li.getSize()).sum();
+            long size = getTotalSize() - sizeBeforePosLedger;
+
+            if (pos.getLedgerId() == currentLedger.getId()) {
+                return size - consumedLedgerSize(currentLedgerSize, 
currentLedgerEntries, pos.getEntryId());
+            } else {
+                return size - consumedLedgerSize(ledgerInfo.getSize(), 
ledgerInfo.getEntries(), pos.getEntryId());
+            }
+        }
+    }
+
     private long consumedLedgerSize(long ledgerSize, long ledgerEntries, long 
consumedEntries) {
         if (ledgerEntries <= 0) {
             return 0;
diff --git 
a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java
 
b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java
index c9021ae22f..ecf6acfc50 100644
--- 
a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java
+++ 
b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java
@@ -277,6 +277,11 @@ public int getTotalNonContiguousDeletedMessagesRange() {
             return 0;
         }
 
+        @Override
+        public long getEstimatedSizeSinceMarkDeletePosition() {
+            return 0L;
+        }
+
         @Override
         public void setThrottleMarkDelete(double throttleMarkDelete) {
         }
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 da15f15f12..a705d156e3 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
@@ -2684,6 +2684,29 @@ public void deleteFailed(ManagedLedgerException 
exception, Object ctx) {
         assertEquals(readPosition, cursor.getReadPosition());
         assertEquals(markDeletePosition, cursor.getMarkDeletedPosition());
     }
-    
+
+    @Test
+    public void testEstimatedUnackedSize() throws Exception {
+        ManagedLedgerConfig config = new ManagedLedgerConfig();
+        config.setMaxEntriesPerLedger(10);
+        ManagedLedger ledger = factory.open("my_test_ledger", new 
ManagedLedgerConfig());
+
+        ManagedCursor cursor = ledger.openCursor("c1");
+
+        byte[] entryData = new byte[5];
+
+        // write 15 entries, saving position of 5th
+        for (int i = 0; i < 4; i++) { ledger.addEntry(entryData); }
+        Position deleteAt = ledger.addEntry(entryData);
+        for (int i = 0; i < 10; i++) { ledger.addEntry(entryData); }
+
+        assertEquals(cursor.getEstimatedSizeSinceMarkDeletePosition(), 15 * 
entryData.length);
+
+        cursor.markDelete(deleteAt);
+
+        // it's not an estimate if all entries are the same size
+        assertEquals(cursor.getEstimatedSizeSinceMarkDeletePosition(), 10 * 
entryData.length);
+    }
+
     private static final Logger log = 
LoggerFactory.getLogger(ManagedCursorTest.class);
 }


 

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
[email protected]


With regards,
Apache Git Services

Reply via email to