merlimat closed pull request #1807: In S3 offloader, dont use 
InputStream#available for stream length
URL: https://github.com/apache/incubator-pulsar/pull/1807
 
 
   

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/OffloadIndexBlock.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/OffloadIndexBlock.java
index 944edcaf09..c7e71c721c 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/OffloadIndexBlock.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/OffloadIndexBlock.java
@@ -19,8 +19,9 @@
 package org.apache.pulsar.broker.s3offload;
 
 import java.io.Closeable;
-import java.io.IOException;
+import java.io.FilterInputStream;
 import java.io.InputStream;
+import java.io.IOException;
 import org.apache.bookkeeper.client.api.LedgerMetadata;
 import org.apache.bookkeeper.common.annotation.InterfaceStability.Unstable;
 
@@ -38,7 +39,7 @@
      *   | index_magic_header | index_block_len | index_entry_count |
      *   | data_object_size | segment_metadata_length | segment metadata | 
index entries ... |
      */
-    InputStream toStream() throws IOException;
+    IndexInputStream toStream() throws IOException;
 
     /**
      * Get the related OffloadIndexEntry that contains the given 
messageEntryId.
@@ -59,9 +60,28 @@
      */
     LedgerMetadata getLedgerMetadata();
 
-    /**
+    /*
      * Get the total size of the data object.
      */
     long getDataObjectLength();
+
+    /**
+     * An input stream which knows the size of the stream upfront.
+     */
+    public static class IndexInputStream extends FilterInputStream {
+        final long streamSize;
+
+        public IndexInputStream(InputStream in, long streamSize) {
+            super(in);
+            this.streamSize = streamSize;
+        }
+
+        /**
+         * @return the number of bytes in the stream.
+         */
+        public long getStreamSize() {
+            return streamSize;
+        }
+    }
 }
 
diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
index ee825323b0..9cb1486bf4 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/S3ManagedLedgerOffloader.java
@@ -180,10 +180,10 @@ static String indexBlockOffloadKey(long ledgerId, UUID 
uuid) {
 
             // upload index block
             try (OffloadIndexBlock index = 
indexBuilder.withDataObjectLength(dataObjectLength).build();
-                 InputStream indexStream = index.toStream()) {
+                 OffloadIndexBlock.IndexInputStream indexStream = 
index.toStream()) {
                 // write the index block
                 ObjectMetadata metadata = new ObjectMetadata();
-                metadata.setContentLength(indexStream.available());
+                metadata.setContentLength(indexStream.getStreamSize());
                 s3client.putObject(new PutObjectRequest(
                     bucket,
                     indexBlockKey,
diff --git 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexBlockImpl.java
 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexBlockImpl.java
index 638edc4451..3c5f33755f 100644
--- 
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexBlockImpl.java
+++ 
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/s3offload/impl/OffloadIndexBlockImpl.java
@@ -28,6 +28,7 @@
 import io.netty.util.Recycler.Handle;
 import java.io.DataInputStream;
 import java.io.IOException;
+import java.io.FilterInputStream;
 import java.io.InputStream;
 import java.util.ArrayList;
 import java.util.List;
@@ -158,15 +159,12 @@ public long getDataObjectLength() {
      *   |segment_metadata_len | segment metadata | index entries |
      */
     @Override
-    public InputStream toStream() throws IOException {
-        int indexBlockLength;
-        int segmentMetadataLength;
+    public OffloadIndexBlock.IndexInputStream toStream() throws IOException {
         int indexEntryCount = this.indexEntries.size();
-
         byte[] ledgerMetadataByte = 
buildLedgerMetadataFormat(this.segmentMetadata);
-        segmentMetadataLength = ledgerMetadataByte.length;
+        int segmentMetadataLength = ledgerMetadataByte.length;
 
-        indexBlockLength = 4 /* magic header */
+        int indexBlockLength = 4 /* magic header */
             + 4 /* index block length */
             + 8 /* data object length */
             + 4 /* segment metadata length */
@@ -191,7 +189,7 @@ public InputStream toStream() throws IOException {
                 .writeInt(entry.getValue().getPartId())
                 .writeLong(entry.getValue().getOffset()));
 
-        return new ByteBufInputStream(out, true);
+        return new OffloadIndexBlock.IndexInputStream(new 
ByteBufInputStream(out, true), indexBlockLength);
     }
 
     static private class InternalLedgerMetadata implements LedgerMetadata {


 

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