This is an automated email from the ASF dual-hosted git repository.

mmerli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git


The following commit(s) were added to refs/heads/master by this push:
     new c01c6be  S3BackedInputStream seeks within current buffer if possible 
(#1892)
c01c6be is described below

commit c01c6beb026f477df2e9874581c24c5bb56b7502
Author: Ivan Kelly <[email protected]>
AuthorDate: Fri Jun 1 20:28:43 2018 +0200

    S3BackedInputStream seeks within current buffer if possible (#1892)
    
    Previously, any time seek was called, the current buffer was
    discarded, and the next time read was called, a request was made to
    S3. This was very inefficient if using seek a lot (which we do).
    
    This patch changes that behaviour so that if the position sought is
    within the range of the current buffer, we just change the reader
    index on the buffer.
    
    Master Issue: #1511
---
 .../s3offload/impl/S3BackedInputStreamImpl.java    | 14 ++++++-
 .../broker/s3offload/S3BackedInputStreamTest.java  | 47 ++++++++++++++++++++++
 2 files changed, 59 insertions(+), 2 deletions(-)

diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/S3BackedInputStreamImpl.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/S3BackedInputStreamImpl.java
index 912a1d5..65f2337 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/S3BackedInputStreamImpl.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/S3BackedInputStreamImpl.java
@@ -47,6 +47,8 @@ public class S3BackedInputStreamImpl extends 
S3BackedInputStream {
     private final int bufferSize;
 
     private long cursor;
+    private long bufferOffsetStart;
+    private long bufferOffsetEnd;
 
     public S3BackedInputStreamImpl(AmazonS3 s3client, String bucket, String 
key,
                                    VersionCheck versionCheck,
@@ -59,6 +61,7 @@ public class S3BackedInputStreamImpl extends 
S3BackedInputStream {
         this.objectLen = objectLen;
         this.bufferSize = bufferSize;
         this.cursor = 0;
+        this.bufferOffsetStart = this.bufferOffsetEnd = -1;
     }
 
     /**
@@ -83,6 +86,8 @@ public class S3BackedInputStreamImpl extends 
S3BackedInputStream {
                 long bytesRead = range[1] - range[0] + 1;
 
                 buffer.clear();
+                bufferOffsetStart = range[0];
+                bufferOffsetEnd = range[1];
                 InputStream s = obj.getObjectContent();
                 int bytesToCopy = (int)bytesRead;
                 while (bytesToCopy > 0) {
@@ -119,8 +124,13 @@ public class S3BackedInputStreamImpl extends 
S3BackedInputStream {
     @Override
     public void seek(long position) {
         log.debug("Seeking to {} on {}/{}, current position {}", position, 
bucket, key, cursor);
-        this.cursor = position;
-        buffer.clear();
+        if (position >= bufferOffsetStart && position <= bufferOffsetEnd) {
+            long newIndex = position - bufferOffsetStart;
+            buffer.readerIndex((int)newIndex);
+        } else {
+            this.cursor = position;
+            buffer.clear();
+        }
     }
 
     @Override
diff --git 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3BackedInputStreamTest.java
 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3BackedInputStreamTest.java
index 9f06ee8..4b75869 100644
--- 
a/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3BackedInputStreamTest.java
+++ 
b/pulsar-broker/src/test/java/org/apache/pulsar/broker/s3offload/S3BackedInputStreamTest.java
@@ -18,6 +18,11 @@
  */
 package org.apache.pulsar.broker.s3offload;
 
+import static org.mockito.Matchers.anyObject;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
 import com.amazonaws.services.s3.AmazonS3;
 import com.amazonaws.services.s3.model.ObjectMetadata;
 
@@ -152,6 +157,48 @@ class S3BackedInputStreamTest extends S3TestBase {
     }
 
     @Test
+    public void testSeekWithinCurrent() throws Exception {
+        String objectKey = "foobar";
+        int objectSize = 12345;
+        RandomInputStream toWrite = new RandomInputStream(0, objectSize);
+
+        ObjectMetadata metadata = new ObjectMetadata();
+        metadata.setContentLength(objectSize);
+        s3client.putObject(BUCKET, objectKey, toWrite, metadata);
+
+        AmazonS3 spiedClient = spy(s3client);
+        S3BackedInputStream toTest = new S3BackedInputStreamImpl(spiedClient, 
BUCKET, objectKey,
+                                                                 (key, md) -> 
{},
+                                                                 objectSize, 
1000);
+
+        // seek forward
+        RandomInputStream firstSeek = new RandomInputStream(0, objectSize);
+        toTest.seek(100);
+        firstSeek.skip(100);
+        for (int i = 0; i < 100; i++) {
+            Assert.assertEquals(firstSeek.read(), toTest.read());
+        }
+
+        // seek forward a bit more, but in same block
+        RandomInputStream secondSeek = new RandomInputStream(0, objectSize);
+        toTest.seek(600);
+        secondSeek.skip(600);
+        for (int i = 0; i < 100; i++) {
+            Assert.assertEquals(secondSeek.read(), toTest.read());
+        }
+
+        // seek back
+        RandomInputStream thirdSeek = new RandomInputStream(0, objectSize);
+        toTest.seek(200);
+        thirdSeek.skip(200);
+        for (int i = 0; i < 100; i++) {
+            Assert.assertEquals(thirdSeek.read(), toTest.read());
+        }
+
+        verify(spiedClient, times(1)).getObject(anyObject());
+    }
+
+    @Test
     public void testSeekForward() throws Exception {
         String objectKey = "foobar";
         int objectSize = 12345;

-- 
To stop receiving notification emails like this one, please contact
[email protected].

Reply via email to