yandrey321 commented on code in PR #11080:
URL: https://github.com/apache/ozone/pull/11080#discussion_r4169417567
##########
hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/DatanodeConfiguration.java:
##########
@@ -855,6 +867,12 @@ public void validate() {
BLOCK_DELETE_COMMAND_WORKER_INTERVAL_DEFAULT;
}
+ if (streamReadFileIdleTimeout.isNegative() ||
streamReadFileIdleTimeout.isZero()) {
Review Comment:
Zero or negative silently resets to the 1m default, so there's no way to
turn the new idle-close off. For a new always-on behavior on the read path, an
operator hitting an unforeseen interaction has no lever except a rebuild.
Suggest treating 0 as disabled (keep the reset for negatives) so pre-patch
behavior is reachable from config.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/StreamBlockInputStream.java:
##########
Review Comment:
DEADLINE_EXCEEDED now routes into handleExceptions, which calls
recordFailedStreamingDatanode() before failing over — so the datanode is added
to failedStreamingDatanodes (line 101), and that set is never cleared for the
life of the stream (it's only ever added to at 521 and read at 213/441).
But a deadline is a property of the call configuration, not of the datanode.
In the scenario the description gives as the motivation for this half of the
patch — a proxy or interceptor imposing its own deadline — every replica will
exceed it identically. So a long-lived reader burns one replica per failover
and after ~replication-factor attempts initStreamRead runs out of candidates
and throws IOException("Failed to start streaming read to any available
DataNodes"), which is not retryable. The patch converts "hangs forever" into
"works for 3 deadlines, then fails hard" rather than into recovery.
Suggest either not recording the exclusion for DEADLINE_EXCEEDED
specifically (it isn't evidence against the peer), or giving the exclusion a
TTL / clearing it when the candidate set is exhausted. Worth also handling here
that an exhausted-candidates failure could fall back to re-including previously
failed datanodes instead of surfacing a hard error.
TestStreamBlockInputStream.java:492 can't catch this: the mock makes
initStreamRead succeed unconditionally on the retry, so the exclusion side
effect is invisible. A version that fails the second datanode would show the
behavior.
##########
hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/GrpcXceiverService.java:
##########
@@ -99,22 +119,63 @@ public StreamObserver<ContainerCommandRequestProto> send(
return new StreamObserver<ContainerCommandRequestProto>() {
private final AtomicBoolean isClosed = new AtomicBoolean(false);
private final RandomAccessFileChannel blockFile = new
RandomAccessFileChannel();
+ // Held while a request is served, so the idle closer never closes the
file under an in-flight read.
+ private final ReentrantLock requestLock = new ReentrantLock();
+ private volatile long lastRequestNanos;
+ private volatile ScheduledFuture<?> idleFileCheck;
boolean close() {
if (isClosed.compareAndSet(false, true)) {
+ final ScheduledFuture<?> check = idleFileCheck;
+ if (check != null) {
+ check.cancel(false);
+ }
blockFile.close();
return true;
}
return false;
}
+ private void scheduleIdleFileCheck(long delayNanos) {
+ try {
+ idleFileCheck = idleFileCloser.schedule(this::closeFileIfIdle,
delayNanos, TimeUnit.NANOSECONDS);
+ } catch (RejectedExecutionException e) {
+ // Server is shutting down; the file is closed when the stream ends.
+ LOG.debug("Idle file closer is shut down, not scheduling check", e);
+ }
+ }
+
+ private void closeFileIfIdle() {
Review Comment:
New failure mode worth acknowledging explicitly, even if it's accepted:
holding the descriptor open gave the stream POSIX unlink immunity — a block
file deleted or replaced underneath an in-progress read kept working. Closing
it on idle gives that up, so a block deletion or container move/replication
landing in the idle window turns the reopen in readBlockImpl into a
FileNotFoundException, surfaced as an IO_EXCEPTION StorageContainerException.
The client treats that as retryable and fails over, which is survivable, but
with finding #1 above each such event also costs a replica. A short comment
here stating that the reopen may legitimately find the file gone, and ideally
mapping a missing file to the same clean error rejectReadBlock produces rather
than a generic internal error, would make the tradeoff deliberate.
##########
hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/XceiverServerGrpc.java:
##########
@@ -234,6 +236,7 @@ public void stop() {
try {
server.shutdown();
server.awaitTermination(5, TimeUnit.SECONDS);
+ xceiverService.shutdown();
Review Comment:
xceiverService.shutdown() is inside the try after
server.awaitTermination(...), so an InterruptedException from awaitTermination
skips it and leaks the ReadBlockIdleFileCloser thread; it's also skipped
entirely on the isStarted == false path. The placement after awaitTermination
is right (in-flight calls may still touch the timer) — just move the call into
a finally. Most visible in MiniOzoneCluster suites that restart datanodes
repeatedly.
##########
hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/GrpcXceiverService.java:
##########
@@ -132,6 +193,11 @@ public void onNext(ContainerCommandRequestProto request) {
if (context != null) {
context.release();
}
+ lastRequestNanos = System.nanoTime();
Review Comment:
lastRequestNanos and the reschedule check run in the finally for every
command type, not just Type.ReadBlock. Two consequences: non-ReadBlock traffic
on the same observer refreshes the idle clock of a block file nothing is
reading, and the legacy ReadChunk path now pays a lock acquire plus a volatile
write per request. Scoping the bookkeeping to the ReadBlock branch keeps the
hot classic path byte-for-byte as it was and makes the invariant ("the clock
tracks reads of this block file") true by construction.
##########
hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/TestGrpcXceiverService.java:
##########
@@ -0,0 +1,165 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hadoop.ozone.container.common.transport.server;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+
+import java.io.File;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.concurrent.atomic.AtomicReference;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID;
+import
org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadBlockRequestProto;
+import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type;
+import org.apache.hadoop.hdds.utils.io.RandomAccessFileChannel;
+import org.apache.hadoop.ozone.container.common.interfaces.ContainerDispatcher;
+import org.apache.ozone.test.GenericTestUtils;
+import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+/**
+ * Tests for the streaming ReadBlock handling in {@link GrpcXceiverService}:
the block file held open for a
+ * stream is closed once the stream is idle and reopened by the next request.
+ */
+class TestGrpcXceiverService {
+
+ private static final Duration IDLE_TIMEOUT = Duration.ofMillis(200);
+
+ @TempDir
+ private Path tempDir;
+
+ private GrpcXceiverService service;
+
+ @AfterEach
+ void shutdown() {
+ if (service != null) {
+ service.shutdown();
+ }
+ }
+
+ private static ContainerCommandRequestProto readBlockRequest() {
+ return ContainerCommandRequestProto.newBuilder()
+ .setCmdType(Type.ReadBlock)
+ .setContainerID(1)
+ .setDatanodeUuid("dn")
+ .setReadBlock(ReadBlockRequestProto.newBuilder()
+
.setBlockID(DatanodeBlockID.newBuilder().setContainerID(1).setLocalID(1))
+ .setOffset(0))
+ .build();
+ }
+
+ /**
+ * Mimics {@code KeyValueHandler.readBlockImpl}: open the block file on
first use, then serve the request.
+ * The optional hook runs while the request is in flight.
+ */
+ private ContainerDispatcher mockDispatcher(File blockFile,
AtomicReference<RandomAccessFileChannel> channelRef,
+ Runnable inFlight) throws Exception {
+ ContainerDispatcher dispatcher = mock(ContainerDispatcher.class);
+ doAnswer(inv -> {
+ RandomAccessFileChannel channel = inv.getArgument(2);
+ channelRef.set(channel);
+ if (!channel.isOpen()) {
+ channel.open(blockFile);
+ }
+ inFlight.run();
+ return null;
+ }).when(dispatcher).streamDataReadOnly(any(), any(), any(), any());
+ return dispatcher;
+ }
+
+ @Test
+ void idleStreamClosesBlockFileAndNextRequestReopensIt() throws Exception {
+ File blockFile = Files.createFile(tempDir.resolve("block")).toFile();
+ AtomicReference<RandomAccessFileChannel> channelRef = new
AtomicReference<>();
+ service = new GrpcXceiverService(mockDispatcher(blockFile, channelRef, ()
-> { }), IDLE_TIMEOUT, "test-");
+ StreamObserver<ContainerCommandResponseProto> responseObserver =
mock(StreamObserver.class);
+ StreamObserver<ContainerCommandRequestProto> requestObserver =
service.send(responseObserver);
+
+ requestObserver.onNext(readBlockRequest());
+ RandomAccessFileChannel channel = channelRef.get();
+ assertNotNull(channel);
+ assertTrue(channel.isOpen(), "block file should be open right after a
request");
+
+ GenericTestUtils.waitFor(() -> !channel.isOpen(), 20, 5000);
+ verify(responseObserver, never()).onError(any());
+ verify(responseObserver, never()).onCompleted();
+
+ requestObserver.onNext(readBlockRequest());
+ assertTrue(channel.isOpen(), "next request should reopen the block file");
Review Comment:
assertTrue(channel.isOpen(), "next request should reopen the block file")
runs after onNext returns, with IDLE_TIMEOUT = 200ms. Any stall longer than
that between the call returning and this line — GC, a loaded CI box — closes
the file and fails the test. Same shape at :105. Either assert openness from
inside the dispatcher hook (where it's guaranteed by the lock) or make the
timeout generous for this case and use the short one only where you're waiting
for a close.
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]