Repository navigation
HDDS-16395. Add lock-free StreamBlock pread - #11314
Conversation
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
One or more issues must be addressed before approval.
Get a fresh assessment by requesting another Copilot review.
Review effort: Lite
Findings: 1
Open (2)
What changed in this PR
Adds lock-free positioned reads for StreamBlockInputStream via one-shot gRPC requests and routes multipart StreamBlock reads through the stateless path.
Changes:
- Added concurrent and EOF positioned-read coverage.
- Refactored streaming readers and multipart routing for stateless reads.
- Added retry and response handling for one-shot reads.
| File | Description |
|---|---|
| hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java | Updated as part of this pull request. |
| hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestMultipartInputStream.java | Updated as part of this pull request. |
| hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java | Updated as part of this pull request. |
| hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/MultipartInputStream.java | Updated as part of this pull request. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
In HDDS-15424 apache#11102 we added stateless positioned reads for the Ratis (BlockInputStream) path only; StreamBlock reads still used synchronized seek-read-restore. Serve StreamBlock preads via one-shot gRPC reads that do not touch the sequential cursor, route them through MultipartInputStream.
1. remove extra readers 2. remove refresh lock 3. adjust other comments
- remove unnecessary comments
peterxcli
left a comment
There was a problem hiding this comment.
Thanks for the update! I've left a few cleanup suggestions below. please feel free to take any that are helpful and ignore the rest.
suggest diff:
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 6ce51edff6..2073d3548f 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
@@ -205,77 +205,50 @@ public int readPositioned(long blockOffset, ByteBuffer dst) throws IOException {
final int startPosition = dst.position();
int callRetries = 0;
while (true) {
- final AtomicReference<StreamingReader> readerRef = new AtomicReference<>();
- try {
- return preadOnce(factory, readerRef, blockOffset, length, dst);
- } catch (IOException e) {
- dst.position(startPosition);
- handlePreadException(e, readerRef.get(), blockOffset, callRetries++);
- }
- }
- }
-
- private int preadOnce(XceiverClientFactory factory, AtomicReference<StreamingReader> readerRef,
- long blockOffset, int length, ByteBuffer dst) throws IOException {
- final Pipeline pipeline = pipelineRef.get();
- final XceiverClientGrpc client = acquireGrpcClientForReadData(factory, pipeline);
- try {
+ final Pipeline pipeline = pipelineRef.get();
+ final XceiverClientGrpc client = acquireGrpcClientForReadData(factory, pipeline);
final StreamingReader reader = new StreamingReader(client);
- readerRef.set(reader);
- boolean streamInitialized = false;
try {
+ // On failure initStreamRead releases its own permit, so the reader is closed only after it succeeds.
client.initStreamRead(blockID, reader, failedStreamingDatanodes);
- streamInitialized = true;
- final StreamingReadResponse response = reader.getResponse();
- if (response == null) {
- throw new IOException("Uninitialized StreamingReadResponse: " + blockID);
- }
- client.streamRead(ContainerProtocolCalls.buildReadBlockCommandProto(
- blockID, blockOffset, length, responseDataSize, tokenRef.get(), pipeline), response);
-
- int copied = 0;
- while (copied < length) {
- final ReadBlockResponseProto proto = reader.poll();
- if (proto == null) {
- break;
- }
- final ByteBuffer buffer = getByteBuffer(proto, blockOffset + copied);
- if (buffer == null || !buffer.hasRemaining()) {
- continue;
- }
- buffer.limit(buffer.position() + Math.min(buffer.remaining(), length - copied));
- copied += buffer.remaining();
- dst.put(buffer);
+ try {
+ return preadOnce(client, reader, pipeline, blockOffset, length, dst);
+ } finally {
+ closeStream(reader);
}
- return copied > 0 ? copied : EOF;
+ } catch (IOException e) {
+ dst.position(startPosition);
+ handlePreadException(e, reader, blockOffset, callRetries++);
} finally {
- closePreadReader(reader, streamInitialized);
+ factory.releaseClientForReadData(client, false);
}
- } finally {
- factory.releaseClientForReadData(client, false);
}
}
- private void closePreadReader(StreamingReader reader, boolean releasePermit) {
- if (releasePermit) {
- reader.onCompleted();
- }
+ private int preadOnce(XceiverClientGrpc client, StreamingReader reader, Pipeline pipeline,
+ long blockOffset, int length, ByteBuffer dst) throws IOException {
final StreamingReadResponse response = reader.getResponse();
if (response == null) {
- return;
+ throw new IOException("Uninitialized StreamingReadResponse: " + blockID);
}
- final ClientCallStreamObserver<ContainerProtos.ContainerCommandRequestProto> requestObserver =
- response.getRequestObserver();
- try {
- requestObserver.onCompleted();
- } catch (RuntimeException e) {
- LOG.warn("Failed to close gRPC request stream for {}", reader, e);
- try {
- requestObserver.cancel(STREAM_CLOSE_REASON, e);
- } catch (RuntimeException cancelEx) {
- LOG.warn("Failed to cancel gRPC request stream for {}", reader, cancelEx);
+ client.streamRead(ContainerProtocolCalls.buildReadBlockCommandProto(
+ blockID, blockOffset, length, responseDataSize, tokenRef.get(), pipeline), response);
+
+ int copied = 0;
+ while (copied < length) {
+ final ReadBlockResponseProto proto = reader.poll();
+ if (proto == null) {
+ break;
}
+ final ByteBuffer buffer = getByteBuffer(proto, blockOffset + copied);
+ if (buffer == null || !buffer.hasRemaining()) {
+ continue;
+ }
+ buffer.limit(buffer.position() + Math.min(buffer.remaining(), length - copied));
+ copied += buffer.remaining();
+ dst.put(buffer);
}
+ return copied > 0 ? copied : EOF;
}
private void handlePreadException(IOException cause, StreamingReader reader, long blockOffset,
@@ -287,27 +260,16 @@ private void handlePreadException(IOException cause, StreamingReader reader, lon
eof.initCause(cause);
throw eof;
}
- final StorageContainerException sce = findStorageContainerException(cause);
- if (sce == null && !isConnectivityIssue(root) && !(root instanceof TimeoutIOException)) {
- throw cause;
- }
- if (!shouldRetryRead(root, retryPolicy, callRetries)) {
+ final boolean retriable = root instanceof StorageContainerException || isConnectivityIssue(root)
+ || root instanceof TimeoutIOException;
+ if (!retriable || !shouldRetryRead(root, retryPolicy, callRetries)) {
throw cause;
}
recordFailedStreamingDatanode(reader);
- refreshBlockInfo(root, blockID, pipelineRef, tokenRef, refreshFunction);
+ refreshBlockInfo(root);
LOG.warn("Refreshing block data to pread block {} due to {}", blockID, cause.getMessage());
}
- private static StorageContainerException findStorageContainerException(Throwable throwable) {
- for (Throwable t = throwable; t != null; t = t.getCause()) {
- if (t instanceof StorageContainerException) {
- return (StorageContainerException) t;
- }
- }
- return null;
- }
-
private synchronized boolean dataAvailableToRead(int length, boolean preRead) throws IOException {
if (position >= blockLength) {
return false;
@@ -415,7 +377,11 @@ private synchronized void closeReader(String reason) {
final StreamingReader reader = streamingReader;
streamingReader = null;
LOG.debug("{} closeReader for {}", getName(reader), reason);
+ closeStream(reader);
+ }
+ /** Release the stream read permit of the given reader and close its gRPC request stream. */
+ private static void closeStream(StreamingReader reader) {
reader.onCompleted();
final StreamingReadResponse response = reader.getResponse();
@@ -461,6 +427,7 @@ private XceiverClientGrpc acquireGrpcClientForReadData(XceiverClientFactory fact
throw new IOException("Failed to acquire client for " + pipeline);
}
if (!(acquired instanceof XceiverClientGrpc)) {
+ factory.releaseClientForReadData(acquired, false);
throw new IOException("Unexpected client class: " + acquired.getClass().getName() + ", " + pipeline);
}
return (XceiverClientGrpc) acquired;
@@ -614,6 +581,7 @@ public String toString() {
*/
public class StreamingReader implements StreamingReaderSpi {
private final String name = StreamBlockInputStream.this.name + "-reader" + READER_ID.getAndIncrement();
+ /** The client whose stream read permit this reader holds. */
private final XceiverClientGrpc client;
/** Response queue: poll is blocking while offer is non-blocking. */
@@ -623,42 +591,10 @@ public class StreamingReader implements StreamingReaderSpi {
private final AtomicBoolean semaphoreReleased = new AtomicBoolean(false);
private final AtomicReference<StreamingReadResponse> response = new AtomicReference<>();
- public StreamingReader(XceiverClientGrpc client) {
+ StreamingReader(XceiverClientGrpc client) {
this.client = client;
}
- private ReadBuffer read(int length, boolean preRead) throws IOException {
- checkError();
- if (isDone()) {
- if (isQueueEmpty()) {
- return null;
- }
- } else {
- readBlock(length, preRead);
- }
-
- while (true) {
- final ReadBlockResponseProto proto = poll();
- if (proto == null) {
- return null;
- }
- 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 {}", this, length, preRead, read);
- return read;
- }
- }
- }
-
- final boolean isDone() {
- return future.isDone();
- }
-
- final boolean isQueueEmpty() {
- return responseQueue.isEmpty();
- }
-
void checkError() throws IOException {
if (future.isCompletedExceptionally()) {
try {
@@ -706,6 +642,34 @@ ReadBlockResponseProto poll() 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.
+ if (responseQueue.isEmpty()) {
+ return null;
+ }
+ } else {
+ // send gRPC onNext(..)
+ readBlock(length, preRead);
+ }
+
+ // poll buffer from queue
+ while (true) {
+ final ReadBlockResponseProto proto = poll();
+ if (proto == null) {
+ return null;
+ }
+ 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;
+ }
+ }
+ }
+
private void releaseResources() {
if (semaphoreReleased.compareAndSet(false, true)) {
client.completeStreamRead();
@@ -779,7 +743,7 @@ public void onCompleted() {
releaseResources();
}
- public StreamingReadResponse getResponse() {
+ StreamingReadResponse getResponse() {
return response.get();
}
diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java
index 585759c241..88258bfe54 100644
--- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java
+++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestStreamBlockInputStream.java
@@ -148,6 +148,8 @@ public void testPositionedReadAtEof() throws Exception {
ByteBuffer dst = ByteBuffer.allocate(10);
assertEquals(-1, sbis.readPositioned(data.length, dst));
assertEquals(-1, sbis.readPositioned(data.length + 1, dst));
+ // An empty buffer also gets -1, as PartInputStream#readPositioned documents.
+ assertEquals(-1, sbis.readPositioned(0, ByteBuffer.allocate(0)));
}
}
@@ -157,17 +159,11 @@ public void testPositionedReadDoesNotReleasePermitWhenInitStreamReadFails() thro
clientConfig.setMaxReadRetryCount(0);
BlockID blockID = new BlockID(1L, 18L);
Pipeline pipeline = mockStandalonePipeline();
- ClientCallStreamObserver<ContainerCommandRequestProto> requestObserver =
- mock(ClientCallStreamObserver.class);
- StreamingReadResponse streamingReadResponse = mock(StreamingReadResponse.class);
- when(streamingReadResponse.getRequestObserver()).thenReturn(requestObserver);
+ // Like XceiverClientGrpc, fail before the gRPC call starts: initStreamRead then releases its own permit.
XceiverClientGrpc xceiverClient = mock(XceiverClientGrpc.class);
- doAnswer(invocation -> {
- StreamingReaderSpi reader = invocation.getArgument(1);
- reader.setStreamingReadResponse(streamingReadResponse);
- throw new IOException("initStreamRead failed");
- }).when(xceiverClient).initStreamRead(any(BlockID.class), any(), any());
+ doThrow(new IOException("initStreamRead failed"))
+ .when(xceiverClient).initStreamRead(any(BlockID.class), any(), any());
XceiverClientFactory xceiverClientFactory = mock(XceiverClientFactory.class);
when(xceiverClientFactory.acquireClientForReadData(any(Pipeline.class)))
@@ -181,7 +177,6 @@ public void testPositionedReadDoesNotReleasePermitWhenInitStreamReadFails() thro
}
verify(xceiverClient, never()).completeStreamRead();
- verify(requestObserver, times(1)).onCompleted();
}
@Test
@@ -402,7 +397,7 @@ private XceiverClientGrpc mockOffsetAwareStreamingReadClient(byte[] data,
}
byte[] slice = Arrays.copyOfRange(data, start, end);
reader.onNext(buildResponseProto(slice, offset));
- if (offset > 0 || end >= data.length) {
+ if (end >= data.length) {
reader.onCompleted();
}
return null;
@@ -963,6 +958,68 @@ public void testFailFastOnGrpcOutOfRangeStatus() throws Exception {
}
}
+ /**
+ * A streamed error response (e.g. CONTAINER_NOT_FOUND) is retried after refreshing the block location,
+ * for both sequential and positioned reads.
+ */
+ @Test
+ public void testRetriesErrorResponseAfterRefresh() throws Exception {
+ OzoneClientConfig clientConfig = newStreamReadConfig();
+ clientConfig.setMaxReadRetryCount(1);
+ clientConfig.setReadRetryInterval(0);
+ byte[] data = new byte[] {1, 2, 3, 4};
+ ClientCallStreamObserver<ContainerCommandRequestProto> requestObserver = mock(ClientCallStreamObserver.class);
+ StreamingReadResponse streamingReadResponse = mock(StreamingReadResponse.class);
+ when(streamingReadResponse.getRequestObserver()).thenReturn(requestObserver);
+ when(streamingReadResponse.getDatanodeDetails()).thenReturn(MockDatanodeDetails.randomDatanodeDetails());
+
+ // Every other ReadBlock request gets an error response; only one reader is active at a time here.
+ AtomicInteger streamReads = new AtomicInteger();
+ AtomicReference<StreamingReaderSpi> lastReader = new AtomicReference<>();
+ XceiverClientGrpc xceiverClient = mock(XceiverClientGrpc.class);
+ doAnswer(inv -> {
+ StreamingReaderSpi reader = inv.getArgument(1);
+ reader.setStreamingReadResponse(streamingReadResponse);
+ lastReader.set(reader);
+ return null;
+ }).when(xceiverClient).initStreamRead(any(BlockID.class), any(), any());
+ doAnswer(inv -> {
+ StreamingReaderSpi reader = lastReader.get();
+ if (streamReads.getAndIncrement() % 2 == 0) {
+ reader.onNext(ContainerCommandResponseProto.newBuilder()
+ .setCmdType(Type.ReadBlock)
+ .setResult(ContainerProtos.Result.CONTAINER_NOT_FOUND)
+ .setMessage("Container not found")
+ .build());
+ } else {
+ reader.onNext(buildResponseProto(data, 0));
+ reader.onCompleted();
+ }
+ return null;
+ }).when(xceiverClient).streamRead(any(), any());
+ XceiverClientFactory xceiverClientFactory = mock(XceiverClientFactory.class);
+ when(xceiverClientFactory.acquireClientForReadData(any(Pipeline.class))).thenReturn(xceiverClient);
+ AtomicInteger refreshes = new AtomicInteger();
+ Function<BlockID, BlockLocationInfo> refresh = b -> {
+ refreshes.incrementAndGet();
+ return null;
+ };
+
+ try (StreamBlockInputStream sbis = new StreamBlockInputStream(
+ new BlockID(1L, 22L), data.length, mockStandalonePipeline(), null, xceiverClientFactory,
+ refresh, clientConfig)) {
+ ByteBuffer positioned = ByteBuffer.allocate(data.length);
+ assertEquals(data.length, sbis.readPositioned(0, positioned));
+ assertArrayEquals(data, positioned.array());
+ assertEquals(1, refreshes.get());
+
+ ByteBuffer sequential = ByteBuffer.allocate(data.length);
+ assertEquals(data.length, sbis.read(sequential));
+ assertArrayEquals(data, sequential.array());
+ assertEquals(2, refreshes.get());
+ }
+ }
+
/**
* Mocks a streaming read client which captures the StreamingReaderSpi during
* initStreamRead and drives the given callback when a ReadBlock request is sent.
diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/ConnectionFailureUtils.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/ConnectionFailureUtils.java
index d7f45b1cd7..9ba515e91a 100644
--- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/ConnectionFailureUtils.java
+++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/ConnectionFailureUtils.java
@@ -120,7 +120,7 @@ public static IOException unwrapCause(IOException ex) {
t = t.getCause();
continue;
}
- if (t.getCause() instanceof IOException) {
+ if (t.getCause() instanceof IOException || t.getCause() instanceof ExecutionException) {
t = t.getCause();
continue;
}
diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestConnectionFailureUtils.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestConnectionFailureUtils.java
index 00c327f689..4c766a2a56 100644
--- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestConnectionFailureUtils.java
+++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestConnectionFailureUtils.java
@@ -18,6 +18,7 @@
package org.apache.hadoop.hdds.utils;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.io.EOFException;
@@ -27,7 +28,10 @@
import java.net.SocketException;
import java.net.SocketTimeoutException;
import java.net.UnknownHostException;
+import java.util.concurrent.ExecutionException;
import java.util.stream.Stream;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result;
+import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
import org.apache.hadoop.security.AccessControlException;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
@@ -117,6 +121,14 @@ public void testNullIsNotAConnectionFailure() {
assertFalse(ConnectionFailureUtils.isConnectionFailure(null));
}
+ /** The shape StreamBlockInputStream's reader throws for a failed stream. */
+ @Test
+ public void testUnwrapCauseThroughIOExceptionWrappingExecutionException() {
+ StorageContainerException sce = new StorageContainerException("not found", Result.CONTAINER_NOT_FOUND);
+ IOException wrapped = new IOException("Streaming read failed", new ExecutionException(sce));
+ assertSame(sce, ConnectionFailureUtils.unwrapCause(wrapped));
+ }
+
/**
* {@code Throwable.initCause} contractually rejects setting cause to
* the throwable itself, but cycles of length 2+ have appeared in| /** | ||
| * Default positioned read: seek-read-restore on this part's cursor. Block streams with a | ||
| * stateless implementation override this. | ||
| */ | ||
| @Override | ||
| public int readPositioned(long partOffset, ByteBuffer dst) throws IOException { |
There was a problem hiding this comment.
Do you want to do this in other PR? because if the EC support positional read, too. we can remove the statelessSupported flag and its logic from:
to sth like:
int i = 0;
long streamLength = 0L;
for (PartInputStream partInputStream : inputStreams) {
this.partOffsets[i++] = streamLength;
streamLength += partInputStream.getLength();
}
this.length = streamLength;There was a problem hiding this comment.
let's keep it for now and does not expect EC to support stateless/lock-free pread at the moment
There was a problem hiding this comment.
@taklwu do you know if we already ticket to track this?
most of the changes make senses, I have applied them into the current PR. |
peterxcli
left a comment
There was a problem hiding this comment.
LGTM, thanks for the update!
|
re-triggered the failed ci job |
|
@taklwu thanks! |


What changes were proposed in this pull request?
HDDS-16395. Add lock-free StreamBlock pread
Please describe your PR in detail:
In HDDS-15424 #11102 we added stateless positioned reads for the Ratis (BlockInputStream) path only; StreamBlock reads still used synchronized seek-read-restore. Serve StreamBlock preads via one-shot gRPC reads that do not touch the sequential cursor, route them through MultipartInputStream.
What is the link to the Apache JIRA
https://issues.apache.org/jira/browse/HDDS-16395
How was this patch tested?
unit tests