joseluisll commented on code in PR #8777:
URL: https://github.com/apache/hadoop/pull/8777#discussion_r4211722396
##########
hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/fsdataset/impl/FsDatasetImpl.java:
##########
@@ -3091,6 +3091,22 @@ static ReplicaRecoveryInfo
initReplicaRecoveryImpl(String bpid, ReplicaMap map,
+ replica);
}
+ // A packet write that fails after its data reached the block file but
+ // before bytesOnDisk was updated (e.g. ClosedByInterruptException while
+ // syncing) leaves an unacknowledged tail in the block file. Drop it, as
+ // recoverRbwImpl does for pipeline recovery (HDFS-11472), instead of
+ // failing checkReplicaFiles and excluding the replica from recovery.
Review Comment:
The comment says this mirrors `recoverRbwImpl`, but that method trusts the
file first (it raises `bytesOnDisk` to `blockDataLength`) and only then
truncates to `bytesAcked`. A DataNode restart also keeps this tail when its
checksums are valid, since the replica length is rebuilt by
`validateIntegrityAndSetLength`.
Truncating to `bytesOnDisk` seems like a fine conservative choice, since the
tail's checksums may never have reached the meta file. Could the comment say
that, instead of citing `recoverRbwImpl`?
##########
hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/fsdataset/impl/FsDatasetImpl.java:
##########
@@ -3091,6 +3091,22 @@ static ReplicaRecoveryInfo
initReplicaRecoveryImpl(String bpid, ReplicaMap map,
+ replica);
}
+ // A packet write that fails after its data reached the block file but
+ // before bytesOnDisk was updated (e.g. ClosedByInterruptException while
+ // syncing) leaves an unacknowledged tail in the block file. Drop it, as
+ // recoverRbwImpl does for pipeline recovery (HDFS-11472), instead of
+ // failing checkReplicaFiles and excluding the replica from recovery.
+ final long bytesOnDisk = replica.getBytesOnDisk();
+ final long blockDataLength = replica.getBlockDataLength();
+ if (blockDataLength > bytesOnDisk) {
Review Comment:
Minor: the truncation runs before the generation-stamp and recovery-id
checks below. Harmless in practice, since a replica that fails them is stale
anyway, but moving it (together with `checkReplicaFiles`) next to the RUR
conversion would keep the on-disk change on the path that succeeds.
##########
hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestLeaseRecovery.java:
##########
@@ -242,6 +249,79 @@ public void testBlockRecoveryWithLessMetafile() throws
Exception {
assertEquals(newFileLen, expectedNewFileLen);
}
+ /**
+ * A packet write that fails after its data reached the block file but
+ * before bytesOnDisk was updated leaves an RBW replica whose block file is
+ * longer than bytesOnDisk. Lease recovery must drop the unacknowledged tail
+ * and recover the replica instead of rejecting it with
+ * "Block length mismatch".
+ */
+ @Test
+ @Timeout(120)
+ public void testLeaseRecoveryAfterInterruptedPacketWrite() throws Exception {
+ Configuration conf = new HdfsConfiguration();
+ // A failed block recovery is only retried after 30 heartbeat intervals,
+ // well after this test stops waiting for the lease to be recovered.
+ conf.setLong(DFSConfigKeys.DFS_HEARTBEAT_INTERVAL_KEY, 1);
Review Comment:
As far as I can tell, a retry wouldn't succeed anyway, because the mismatch
is held in DataNode memory, and with the default 3s heartbeat the retry would
come at 90s, also past the 15s wait. The 1s heartbeat seems to be about getting
the recovery command delivered quickly. Could the comment say that?
Also, without the fix the test fails with a bare `TimeoutException`; a
message on the `waitFor` would make that failure easier to read.
##########
hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestLeaseRecovery.java:
##########
@@ -242,6 +249,79 @@ public void testBlockRecoveryWithLessMetafile() throws
Exception {
assertEquals(newFileLen, expectedNewFileLen);
}
+ /**
+ * A packet write that fails after its data reached the block file but
+ * before bytesOnDisk was updated leaves an RBW replica whose block file is
+ * longer than bytesOnDisk. Lease recovery must drop the unacknowledged tail
+ * and recover the replica instead of rejecting it with
+ * "Block length mismatch".
+ */
+ @Test
+ @Timeout(120)
+ public void testLeaseRecoveryAfterInterruptedPacketWrite() throws Exception {
Review Comment:
This end-to-end test depends on `delayWriteToDisk` staying between
`writeDataToDisk` and `flushOrSync` (the `assertFalse` and block-length
assertions guard that, which is nice). Would a direct `initReplicaRecovery`
test be worth adding as well, in the style of
`TestInterDatanodeProtocol#testUpdateReplicaUnderRecovery`? It could cover an
RBW replica with extra bytes in the block file, with and without matching
checksums in the meta file, and the `bytesOnDisk == 0` case (first packet
interrupted), where the replica is truncated to 0 and recovery deletes the
block.
##########
hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestLeaseRecovery.java:
##########
@@ -242,6 +249,79 @@ public void testBlockRecoveryWithLessMetafile() throws
Exception {
assertEquals(newFileLen, expectedNewFileLen);
}
+ /**
+ * A packet write that fails after its data reached the block file but
+ * before bytesOnDisk was updated leaves an RBW replica whose block file is
+ * longer than bytesOnDisk. Lease recovery must drop the unacknowledged tail
+ * and recover the replica instead of rejecting it with
+ * "Block length mismatch".
+ */
+ @Test
+ @Timeout(120)
+ public void testLeaseRecoveryAfterInterruptedPacketWrite() throws Exception {
+ Configuration conf = new HdfsConfiguration();
+ // A failed block recovery is only retried after 30 heartbeat intervals,
+ // well after this test stops waiting for the lease to be recovered.
+ conf.setLong(DFSConfigKeys.DFS_HEARTBEAT_INTERVAL_KEY, 1);
+ cluster = new MiniDFSCluster.Builder(conf).numDataNodes(1).build();
+ cluster.waitActive();
+ DistributedFileSystem dfs = cluster.getFileSystem();
+ Path file = new Path("/testLeaseRecoveryAfterInterruptedPacketWrite");
+ byte[] acked = AppendTestUtil.randomBytes(0xFEEDL, 4096);
+ byte[] unacked = AppendTestUtil.randomBytes(0xBEEFL, 1024);
+
+ FSDataOutputStream out = dfs.create(file, (short) 1);
+ out.write(acked);
+ out.hsync();
+
+ AtomicBoolean interruptNextWrite = new AtomicBoolean(true);
+ DataNodeFaultInjector oldInjector = DataNodeFaultInjector.get();
+ DataNodeFaultInjector.set(new DataNodeFaultInjector() {
+ @Override
+ public void delayWriteToDisk() {
+ // The packet data is in the block file. Interrupt the receiver, as the
+ // PacketResponder does on an ack failure, so the following fsync fails
+ // with ClosedByInterruptException before bytesOnDisk is updated.
+ if (interruptNextWrite.getAndSet(false)) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ });
+ try {
+ out.write(unacked);
+ out.hsync();
+ fail("hsync should fail after the interrupted packet write");
+ } catch (IOException expected) {
+ // the only datanode in the pipeline failed the write
+ } finally {
+ DataNodeFaultInjector.set(oldInjector);
+ }
+ assertFalse(interruptNextWrite.get(), "packet write was not interrupted");
+ ((DFSOutputStream) out.getWrappedStream()).abort();
+
+ ExtendedBlock block = cluster.getNameNodeRpc()
+ .getBlockLocations(file.toString(), 0, Long.MAX_VALUE).get(0)
+ .getBlock();
+ ReplicaInfo rbw = FsDatasetTestUtil.fetchReplicaInfo(
+ DataNodeTestUtils.getFSDataset(cluster.getDataNodes().get(0)),
+ block.getBlockPoolId(), block.getBlockId());
+ assertEquals(acked.length, rbw.getBytesOnDisk());
+ assertEquals(acked.length + unacked.length, rbw.getBlockDataLength(),
+ "block file should contain the unacknowledged packet");
+
+ DistributedFileSystem newDfs = (DistributedFileSystem) FileSystem
+ .newInstance(cluster.getConfiguration(0));
+ GenericTestUtils.waitFor(() -> {
+ try {
+ return newDfs.recoverLease(file);
+ } catch (IOException e) {
+ return false;
+ }
+ }, 500, 15000);
+ assertEquals(acked.length, newDfs.getFileStatus(file).getLen());
+ assertArrayEquals(acked, DFSTestUtil.readFileAsBytes(newDfs, file));
Review Comment:
Could you also assert on the DataNode side after recovery, e.g. that the
replica is FINALIZED and its block file length equals `acked.length`? That
would check the truncation directly rather than only through the client read.
--
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]