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]