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