merlimat closed pull request #1892: S3BackedInputStream seeks within current 
buffer if possible
URL: https://github.com/apache/incubator-pulsar/pull/1892
 
 
   

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/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 912a1d514b..65f233783e 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 @@
     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 S3BackedInputStreamImpl(AmazonS3 s3client, String 
bucket, String key,
         this.objectLen = objectLen;
         this.bufferSize = bufferSize;
         this.cursor = 0;
+        this.bufferOffsetStart = this.bufferOffsetEnd = -1;
     }
 
     /**
@@ -83,6 +86,8 @@ private boolean refillBufferIfNeeded() throws IOException {
                 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 int read(byte[] b, int off, int len) throws 
IOException {
     @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 9f06ee8674..4b758695e4 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;
 
@@ -151,6 +156,48 @@ public void testSeek() throws Exception {
         }
     }
 
+    @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";


 

----------------------------------------------------------------
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