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