This is an automated email from the ASF dual-hosted git repository.

szetszwo pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/master by this push:
     new 50978e3ad7c HDDS-15422. Stream read seek should not close stream 
(#10415)
50978e3ad7c is described below

commit 50978e3ad7c55e3256582d55b0b9ae9165a911ca
Author: Tsz-Wo Nicholas Sze <[email protected]>
AuthorDate: Fri Aug 14 09:47:00 2026 -0700

    HDDS-15422. Stream read seek should not close stream (#10415)
---
 .../hdds/scm/storage/MultipartInputStream.java     |   4 +-
 .../hdds/scm/storage/StreamBlockInputStream.java   | 186 +++++++++++++--------
 .../rpc/read/TestStreamBlockInputStream.java       |  56 +++++--
 .../ozone/client/rpc/read/TestStreamRead.java      |   1 +
 4 files changed, 163 insertions(+), 84 deletions(-)

diff --git 
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/MultipartInputStream.java
 
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/MultipartInputStream.java
index 075ab08fe5a..e7c259ce464 100644
--- 
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/MultipartInputStream.java
+++ 
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/MultipartInputStream.java
@@ -174,7 +174,9 @@ public synchronized void seek(long pos) throws IOException {
     }
 
     // Reset the previous partStream's position
-    partStreams.get(prevPartIndex).seek(0);
+    if (prevPartIndex != partIndex) {
+      partStreams.get(prevPartIndex).seek(0);
+    }
 
     // Reset all the partStreams above the partIndex. We do this to reset
     // any previous reads which might have updated the higher part
diff --git 
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java
 
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java
index c15cd338908..b99d88ad305 100644
--- 
a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java
+++ 
b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java
@@ -87,7 +87,7 @@ public class StreamBlockInputStream extends 
BlockExtendedInputStream {
   private XceiverClientFactory xceiverClientFactory;
   private XceiverClientGrpc xceiverClient;
 
-  private ByteBuffer buffer;
+  private ReadBuffer readBuffer;
   private long position = 0;
   private long requestedLength = 0;
   private StreamingReader streamingReader;
@@ -116,6 +116,8 @@ public StreamBlockInputStream(
     this.responseDataSize = config.getStreamReadResponseDataSize();
     this.readTimeout = config.getStreamReadTimeout();
     this.readTimeoutNanos = readTimeout.toNanos();
+
+    LOG.debug("{}: new StreamBlockInputStream", name);
   }
 
   @Override
@@ -136,11 +138,12 @@ public synchronized long getPos() {
   @Override
   public synchronized int read() throws IOException {
     checkOpen();
-    if (!dataAvailableToRead(1, true)) {
+    final boolean preRead = true;
+    if (!dataAvailableToRead(1, preRead)) {
       return EOF;
     }
-    int value = buffer.get();
-    advancePosition(1);
+    final int value = readBuffer.getByteBuffer().get();
+    advancePosition(1, preRead);
     return value;
   }
 
@@ -162,12 +165,14 @@ synchronized int readFully(ByteBuffer targetBuf, boolean 
preRead) throws IOExcep
       if (!dataAvailableToRead(targetBuf.remaining(), preRead)) {
         break;
       }
+
+      final ByteBuffer buffer = readBuffer.getByteBuffer();
       int toCopy = Math.min(buffer.remaining(), targetBuf.remaining());
       ByteBuffer tmpBuf = buffer.duplicate();
       tmpBuf.limit(tmpBuf.position() + toCopy);
       targetBuf.put(tmpBuf);
       buffer.position(tmpBuf.position());
-      advancePosition(toCopy);
+      advancePosition(toCopy, preRead);
       read += toCopy;
     }
     return read > 0 ? read : EOF;
@@ -177,30 +182,31 @@ private synchronized boolean dataAvailableToRead(int 
length, boolean preRead) th
     if (position >= blockLength) {
       return false;
     }
+
     while (true) {
       try {
         initialize();
-        if (bufferHasRemaining()) {
-          return true;
+        if (!hasRemaining(readBuffer)) {
+          readBuffer = streamingReader.read(length, preRead);
         }
-        buffer = streamingReader.read(length, preRead);
         retries = 0;
-        return bufferHasRemaining();
+        return hasRemaining(readBuffer);
       } catch (IOException ex) {
         handleExceptions(ex);
       }
     }
   }
 
-  private synchronized void advancePosition(long delta) {
+  private synchronized void advancePosition(long delta, boolean preRead) {
+    LOG.trace("{}: advance {} -> {}", getName(streamingReader), position, 
position + delta);
     position += delta;
-    if (position >= blockLength && streamingReader != null) {
-      closeStream();
+    if (preRead && position >= blockLength) {
+      closeReader("advancePosition");
     }
   }
 
-  private synchronized boolean bufferHasRemaining() {
-    return buffer != null && buffer.hasRemaining();
+  private static boolean hasRemaining(ReadBuffer read) {
+    return read != null && read.getByteBuffer().hasRemaining();
   }
 
   @Override
@@ -220,10 +226,42 @@ public synchronized void seek(long pos) throws 
IOException {
     if (pos == position) {
       return;
     }
-    LOG.debug("{}: seek {} -> {}", this, position, pos);
-    closeStream();
+    LOG.debug("{}: seek {} -> {}", getName(streamingReader), position, pos);
+    readBuffer = reuseReadBuffer(readBuffer, pos);
     position = pos;
-    requestedLength = pos;
+    if (readBuffer == null) {
+      // Only rewind the request high-watermark when the buffered (already 
requested/served) data cannot be reused;
+      // otherwise we would re-request data that is still buffered.
+      requestedLength = pos;
+    }
+  }
+
+  static ReadBuffer reuseReadBuffer(ReadBuffer previous, long blockOffset) {
+    if (previous != null) {
+      final ByteBuffer buffer = getByteBuffer(previous.getProto(), 
blockOffset);
+      if (buffer != null && buffer.hasRemaining()) {
+        previous.getByteBuffer().position(buffer.position());
+        Preconditions.assertSame(buffer.remaining(), 
previous.getByteBuffer().remaining(), "remaining");
+        return previous;
+      }
+    }
+    return null;
+  }
+
+  static ByteBuffer getByteBuffer(ReadBlockResponseProto proto, long 
blockOffset) {
+    final ByteBuffer buffer = proto.getData().asReadOnlyByteBuffer();
+    // Adjust buffer position since the server always returns data starting at 
checksum boundary.
+    final long protoOffset = proto.getOffset();
+    if (blockOffset < protoOffset) {
+      // This can happen after seek, just drop it for now
+      // TODO: consider to cache the proto, which will be useful when seeking 
back.
+      return null;
+    }
+    final long offset = blockOffset - protoOffset;
+    if (offset > 0) {
+      buffer.position(Math.toIntExact(Math.min(offset, buffer.limit())));
+    }
+    return buffer;
   }
 
   @Override
@@ -238,19 +276,15 @@ public synchronized void unbuffer() {
     releaseClient();
   }
 
-  private synchronized void closeStream() {
+  private synchronized void closeReader(String reason) {
+    readBuffer = null;
     if (streamingReader == null) {
-      buffer = null;
       return;
     }
 
     final StreamingReader reader = streamingReader;
     streamingReader = null;
-    buffer = null;
-
-    if (LOG.isDebugEnabled()) {
-      LOG.debug("Closing {}", reader);
-    }
+    LOG.debug("{} closeReader for {}", getName(reader), reason);
 
     reader.onCompleted();
 
@@ -305,6 +339,7 @@ private synchronized void initialize() throws IOException {
       try {
         acquireClient();
         final StreamingReader reader = new StreamingReader();
+        LOG.debug("{}: new StreamingReader", getName(reader));
         xceiverClient.initStreamRead(blockID, reader, 
failedStreamingDatanodes);
         streamingReader = reader;
       } catch (IOException ioe) {
@@ -373,7 +408,7 @@ private void recordFailedStreamingDatanode() {
 
   protected synchronized void releaseClient() {
     if (xceiverClientFactory != null && xceiverClient != null) {
-      closeStream();
+      closeReader("releaseClient");
       xceiverClientFactory.releaseClientForReadData(xceiverClient, false);
       xceiverClient = null;
     }
@@ -413,6 +448,35 @@ public Duration getReadTimeout() {
     return readTimeout;
   }
 
+  private Object getName(StreamingReader reader) {
+    return reader != null ? reader : name;
+  }
+
+  static class ReadBuffer {
+    private final ReadBlockResponseProto proto;
+    private final ByteBuffer buffer;
+
+    ReadBuffer(ReadBlockResponseProto proto, ByteBuffer buffer) {
+      this.proto = proto;
+      this.buffer = buffer;
+    }
+
+    ReadBlockResponseProto getProto() {
+      return proto;
+    }
+
+    ByteBuffer getByteBuffer() {
+      return buffer;
+    }
+
+    @Override
+    public String toString() {
+      return "ReadBuffer: offset=" + proto.getOffset()
+          + ", dataSize=" + proto.getData().size()
+          + ", buffer=" + buffer;
+    }
+  }
+
   /**
    * Implementation of a StreamObserver used to received and buffer streaming 
GRPC reads.
    */
@@ -462,61 +526,43 @@ ReadBlockResponseProto poll() throws IOException {
         }
 
         final long elapsedNanos = System.nanoTime() - startTime;
-        if (elapsedNanos >= readTimeoutNanos) {
-          setFailedAndThrow(new TimeoutIOException(
-              "Timed out waiting for response after " + readTimeout));
+        if (elapsedNanos >= readTimeoutNanos && !future.isDone()) {
+          final TimeoutIOException e = new TimeoutIOException(
+              this + ": Failed to poll a response, timed out " + readTimeout);
+          if (setFailed(e)) {
+            throw e;
+          }
           return null;
         }
       }
     }
 
-    private ByteBuffer read(int length, boolean preRead) throws IOException {
+    private ReadBuffer read(int length, boolean preRead) throws IOException {
       checkError();
       if (future.isDone()) {
         // Don't return null while items remain in the queue. onNext() may 
have delivered items just before
         // onCompleted() fired.
-        return responseQueue.isEmpty() ? null : readFromQueue();
+        if (responseQueue.isEmpty()) {
+          return null;
+        }
+      } else {
+        // send gRPC onNext(..)
+        readBlock(length, preRead);
       }
 
-      readBlock(length, preRead);
-
+      // poll buffer from queue
       while (true) {
-        final ByteBuffer buf = readFromQueue();
-        if (buf == null) {
-          return null; // Stream ended
+        final ReadBlockResponseProto proto = poll();
+        if (proto == null) {
+          return null;
         }
-        if (buf.hasRemaining()) {
-          return buf;
+        final ByteBuffer buffer = getByteBuffer(proto, getPos());
+        final ReadBuffer read = buffer != null ? new ReadBuffer(proto, buffer) 
: null;
+        if (hasRemaining(read)) {
+          LOG.debug("{}: read(length={}, preRead={}) returns {}", name, 
length, preRead, read);
+          return read;
         }
-        // buf is empty: the server aligned its response to a checksum boundary
-        // before our current position and all bytes were skipped. Fetch the 
next
-        // response, which should start at or after our position.
-      }
-    }
-
-    ByteBuffer readFromQueue() throws IOException {
-      final ReadBlockResponseProto readBlock = poll();
-      if (readBlock == null) {
-        return null; // Stream ended
       }
-      // The server always returns data starting from the last checksum 
boundary. Therefore if the reader position is
-      // ahead of the position we received from the server, we need to adjust 
the buffer position accordingly.
-      final ByteString data = readBlock.getData();
-      final ByteBuffer dataBuffer = data.asReadOnlyByteBuffer();
-      final long blockOffset = readBlock.getOffset();
-      final long pos = getPos();
-      if (pos < blockOffset) {
-        // This should not happen, and if it does, we have a bug.
-        setFailedAndThrow(new IllegalStateException(
-            this + ": out of order, position " + pos + " < block offset " + 
blockOffset));
-      }
-      final long offset = pos - blockOffset;
-      if (offset > 0) {
-        dataBuffer.position(Math.toIntExact(Math.min(offset, 
dataBuffer.limit())));
-      }
-      LOG.debug("{}: return response positon {}, length {} (block offset {}, 
length {})",
-          name, pos, dataBuffer.remaining(), blockOffset, data.size());
-      return dataBuffer;
     }
 
     private void releaseResources() {
@@ -575,12 +621,6 @@ StreamingReadResponse getResponse() {
       return response.get();
     }
 
-    private <T extends Throwable> void setFailedAndThrow(T throwable) throws T 
{
-      if (setFailed(throwable)) {
-        throw throwable;
-      }
-    }
-
     private boolean setFailed(Throwable throwable) {
       final boolean completed = future.completeExceptionally(throwable);
       if (!completed) {
@@ -608,9 +648,9 @@ private void setCompleted() {
     }
 
     private void offerToQueue(ReadBlockResponseProto item) {
-      if (LOG.isDebugEnabled()) {
+      if (LOG.isTraceEnabled()) {
         final ContainerProtos.ChecksumData checksumData = 
item.getChecksumData();
-        LOG.debug("{}: enqueue response offset {}, length {}, numChecksums {}, 
bytesPerChecksum={}",
+        LOG.trace("{}: enqueue response offset {}, length {}, numChecksums {}, 
bytesPerChecksum={}",
             name, item.getOffset(), item.getData().size(),
             checksumData.getChecksumsList().size(), 
checksumData.getBytesPerChecksum());
       }
diff --git 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
index 97b802e78b1..2d4ce5095aa 100644
--- 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
+++ 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamBlockInputStream.java
@@ -23,6 +23,7 @@
 
 import java.io.IOException;
 import java.nio.ByteBuffer;
+import java.util.Random;
 import java.util.concurrent.ThreadLocalRandom;
 import org.apache.hadoop.hdds.conf.OzoneConfiguration;
 import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
@@ -61,7 +62,7 @@ public class TestStreamBlockInputStream extends 
InputStreamTests {
     
GenericTestUtils.setLogLevel(LoggerFactory.getLogger("SCMHATransactionMonitor"),
 Level.ERROR);
     GenericTestUtils.setLogLevel(GrpcXceiverService.class, Level.ERROR);
 
-//    
GenericTestUtils.setLogLevel(LoggerFactory.getLogger(StreamBlockInputStream.class),
 Level.TRACE);
+//    GenericTestUtils.setLogLevel(StreamBlockInputStream.class, Level.DEBUG);
 //    
GenericTestUtils.setLogLevel(LoggerFactory.getLogger(XceiverClientGrpc.class), 
Level.TRACE);
   }
 
@@ -81,7 +82,7 @@ void testReadKey() throws Exception {
       OzoneConfiguration conf = cluster.getConf();
 
       runTestReadKey(DATA_LENGTH, false, conf);
-      for (int i = 0; i < 3; i++) {
+      for (int i = 0; i < 2; i++) {
         final int keyLength = DATA_LENGTH + 
ThreadLocalRandom.current().nextInt(DATA_LENGTH);
         runTestReadKey(keyLength, true, conf);
       }
@@ -185,13 +186,25 @@ void assertData(int pos, int length, ByteBuffer buffer) {
   }
 
   @Test
-  void testAll() throws Exception {
+  void testAllWithPreRead() throws Exception {
+    runTestAll(true);
+  }
+
+  @Test
+  void testAllWithoutPreRead() throws Exception {
+    runTestAll(false);
+  }
+
+  void runTestAll(boolean preRead) throws Exception {
     try (MiniOzoneCluster cluster = newCluster()) {
       cluster.waitForClusterToBeReady();
 
       OzoneConfiguration conf = cluster.getConf();
       OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class);
       clientConfig.setStreamReadBlock(true);
+      if (!preRead) {
+        clientConfig.setStreamReadPreReadSize(0);
+      }
       OzoneConfiguration copy = new OzoneConfiguration(conf);
       copy.setFromObject(clientConfig);
       String keyName = getNewKeyName();
@@ -250,15 +263,38 @@ private void testReadKeyFully(String key) throws 
Exception {
     }
   }
 
+  void assertSeekRead(KeyInputStream in, int position) throws IOException {
+    in.seek(position);
+    int b = in.read();
+    assertEquals(inputData[position], (byte) b, "Read data is not same as 
written data at index " + position);
+  }
+
+  private void runTestSeek(KeyInputStream in, int seekSize, Random random) 
throws IOException {
+    LOG.info("runTestSeek: seekSize={}", seekSize);
+    for (int i = 0; i < 100; i++) {
+      int position = random.nextInt(seekSize);
+      assertSeekRead(in, position);
+    }
+
+    for (int position = 0; position < DATA_LENGTH; position += 
random.nextInt(seekSize)) {
+      assertSeekRead(in, position);
+    }
+
+    for (int position = DATA_LENGTH - 1; position >= 0; position -= 
random.nextInt(seekSize)) {
+      assertSeekRead(in, position);
+    }
+    assertSeekRead(in, 0);
+  }
+
   private void testSeek(String key) throws IOException {
-    java.util.Random random = new java.util.Random();
+    final Random random = new Random();
     try (KeyInputStream keyInputStream = bucket.getKeyInputStream(key)) {
-      for (int i = 0; i < 100; i++) {
-        int position = random.nextInt(DATA_LENGTH);
-        keyInputStream.seek(position);
-        int b = keyInputStream.read();
-        assertEquals(inputData[position], (byte) b, "Read data is not same as 
written data at index " + position);
-      }
+      runTestSeek(keyInputStream, CHUNK_SIZE / 8, random);
+      runTestSeek(keyInputStream, CHUNK_SIZE, random);
+      runTestSeek(keyInputStream, BLOCK_SIZE, random);
+      runTestSeek(keyInputStream, DATA_LENGTH, random);
+
+      // error cases
       StreamBlockInputStream blockStream = (StreamBlockInputStream) 
keyInputStream.getPartStreams().get(0);
       long length = blockStream.getLength();
       blockStream.seek(10);
diff --git 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java
 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java
index ee56eb38825..1838b0d928b 100644
--- 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java
+++ 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/read/TestStreamRead.java
@@ -74,6 +74,7 @@ public class TestStreamRead {
     
GenericTestUtils.setLogLevel(LoggerFactory.getLogger("ExpiredContainerReplicaOpScrubber"),
 Level.ERROR);
     
GenericTestUtils.setLogLevel(LoggerFactory.getLogger("SCMHATransactionMonitor"),
 Level.ERROR);
     GenericTestUtils.setLogLevel(LoggerFactory.getLogger(CodecBuffer.class), 
Level.ERROR);
+//    GenericTestUtils.setLogLevel(StreamBlockInputStream.class, Level.DEBUG);
   }
 
   static final int CHUNK_SIZE = 1 << 20;          // 1MB


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to