szetszwo commented on code in PR #1488:
URL: https://github.com/apache/ratis/pull/1488#discussion_r3471704347


##########
ratis-examples/src/main/java/org/apache/ratis/examples/filestore/FileStoreClient.java:
##########
@@ -197,6 +202,49 @@ public DataStreamOutput getStreamOutput(String path, long 
dataSize, RoutingTable
     return 
client.getDataStreamApi().stream(request.toByteString().asReadOnlyByteBuffer(), 
routingTable);
   }
 
+  public DataStreamInput getStreamInput(String path, long offset, long length) 
{
+    final ReadRequestProto read = ReadRequestProto.newBuilder()
+        .setPath(ProtoUtils.toByteString(path))
+        .setOffset(offset)
+        .setLength(length)
+        .build();
+    return 
client.getDataStreamApi().streamReadOnly(read.toByteString().asReadOnlyByteBuffer());
+  }
+
+  /**
+   * Read file data using streaming read and write it to the given channel.
+   *
+   * @return total number of bytes read.
+   */
+  public long streamRead(String path, long offset, long length, 
WritableByteChannel channel)
+      throws IOException {
+    long total = 0;
+    try (DataStreamInput in = getStreamInput(path, offset, length)) {
+      while (true) {
+        final ReferenceCountedObject<DataStreamReply> ref = 
in.readAsync().join();

Review Comment:
   We should not call join().  Otherwise, it becomes sync'ed.  I think it is 
fine for now and we can improve it later.



##########
ratis-examples/src/main/java/org/apache/ratis/examples/filestore/FileInfo.java:
##########
@@ -88,6 +89,34 @@ ByteString read(CheckedFunction<Path, Path, IOException> 
resolver, long offset,
     }
   }
 
+  void streamRead(CheckedFunction<Path, Path, IOException> resolver, long 
offset, long length,
+      WritableByteChannel stream) throws IOException {
+    if (offset + length > getWriteSize()) {
+      throw new IOException("Failed to read Wrote: offset (=" + offset
+          + " + length (=" + length + ") > size = " + getWriteSize()
+          + ", path=" + getRelativePath());
+    }
+
+    try (SeekableByteChannel in = Files.newByteChannel(
+        resolver.apply(getRelativePath()), StandardOpenOption.READ)) {
+      in.position(offset);
+      long remaining = length;
+      while (remaining > 0) {
+        final int chunkSize = FileStoreCommon.getChunkSize(remaining);
+        final ByteBuffer buffer = ByteBuffer.allocateDirect(chunkSize);
+        final int n = in.read(buffer);
+        if (n <= 0) {
+          break;
+        }
+        buffer.flip();
+        stream.write(buffer);
+        remaining -= n;
+      }
+    } finally {
+      stream.close();
+    }

Review Comment:
   Use 
https://docs.oracle.com/javase/8/docs/api/java/nio/channels/FileChannel.html#transferTo-long-long-java.nio.channels.WritableByteChannel-



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to