From 5adc3f02e537399ee8701985f200bdfcbc95cdfc Mon Sep 17 00:00:00 2001
From: f64116045
Date: Thu, 3 Sep 2026 22:30:40 +0800
Subject: [PATCH 1/2] HDDS-16012. Detect replicas with equal BCSID and
mismatched data checksums
---
.../scm/container/ContainerHealthState.java | 7 +
.../ContainerReplicaChecksumMismatch.java | 68 +++++++
.../container/ContainerReplicaInfoResult.java | 44 +++++
.../container/ReplicationManagerReport.java | 9 +
.../container/TestContainerHealthState.java | 13 +-
.../TestContainerReplicaChecksumMismatch.java | 117 +++++++++++
.../TestReplicationManagerReport.java | 16 ++
.../hadoop/hdds/scm/client/ScmClient.java | 12 ++
.../StorageContainerLocationProtocol.java | 12 ++
...ocationProtocolClientSideTranslatorPB.java | 10 +-
.../src/main/proto/ScmAdminProtocol.proto | 1 +
.../src/main/resources/proto.lock | 8 +-
.../replication/ReplicationManager.java | 18 +-
.../DataChecksumMismatchCheckHandler.java | 131 +++++++++++++
...ocationProtocolServerSideTranslatorPB.java | 6 +-
.../scm/server/SCMClientProtocolServer.java | 11 ++
.../replication/TestReplicationManager.java | 80 ++++++++
.../TestDataChecksumMismatchCheckHandler.java | 181 ++++++++++++++++++
.../scm/cli/ContainerOperationClient.java | 19 +-
.../scm/cli/container/InfoSubcommand.java | 36 +++-
.../scm/cli/container/TestInfoSubCommand.java | 65 ++++++-
.../recon/fsck/ReconReplicationManager.java | 53 ++---
.../fsck/TestReconReplicationManager.java | 5 +
23 files changed, 856 insertions(+), 66 deletions(-)
create mode 100644 hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaChecksumMismatch.java
create mode 100644 hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaInfoResult.java
create mode 100644 hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReplicaChecksumMismatch.java
create mode 100644 hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/health/DataChecksumMismatchCheckHandler.java
create mode 100644 hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/health/TestDataChecksumMismatchCheckHandler.java
diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerHealthState.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerHealthState.java
index 56dd1c5620b6..cb5390de8deb 100644
--- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerHealthState.java
+++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerHealthState.java
@@ -107,6 +107,13 @@ public enum ContainerHealthState {
"Containers in OPEN state without any healthy Pipeline",
"OpenContainersWithoutPipeline"),
+ /**
+ * Replicas with the same BCSID have different data checksums.
+ */
+ DATA_CHECKSUM_MISMATCH((short) 10,
+ "Containers with replicas reporting the same BCSID but different data checksums",
+ "DataChecksumMismatchContainers"),
+
// ========== Actual Combinations Found in Code (100+) ==========
/**
diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaChecksumMismatch.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaChecksumMismatch.java
new file mode 100644
index 000000000000..88855c7af693
--- /dev/null
+++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaChecksumMismatch.java
@@ -0,0 +1,68 @@
+/*
+ * 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.hdds.scm.container;
+
+import java.util.Collection;
+import java.util.Objects;
+import java.util.function.Function;
+import java.util.function.ToLongFunction;
+
+/**
+ * Detects different data checksums reported for the same container replica
+ * sequence ID.
+ */
+public final class ContainerReplicaChecksumMismatch {
+
+ private ContainerReplicaChecksumMismatch() {
+ }
+
+ /**
+ * Returns true when replicas with the same sequence ID report different
+ * non-zero data checksums. If any replica has not reported a sequence ID or
+ * checksum yet, no comparison is made.
+ */
+ public static boolean hasMismatch(Collection replicas,
+ Function sequenceId, ToLongFunction dataChecksum) {
+ Objects.requireNonNull(sequenceId, "sequenceId == null");
+ Objects.requireNonNull(dataChecksum, "dataChecksum == null");
+
+ if (replicas == null || replicas.size() < 2) {
+ return false;
+ }
+
+ for (T replica : replicas) {
+ Long replicaSequenceId = sequenceId.apply(replica);
+ long replicaDataChecksum = dataChecksum.applyAsLong(replica);
+ if (replicaSequenceId == null || replicaDataChecksum == 0) {
+ return false;
+ }
+ }
+
+ for (T left : replicas) {
+ Long leftSequenceId = sequenceId.apply(left);
+ long leftDataChecksum = dataChecksum.applyAsLong(left);
+ for (T right : replicas) {
+ if (leftSequenceId.equals(sequenceId.apply(right)) &&
+ leftDataChecksum != dataChecksum.applyAsLong(right)) {
+ return true;
+ }
+ }
+ }
+ return false;
+ }
+}
diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaInfoResult.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaInfoResult.java
new file mode 100644
index 000000000000..cc2b2e03b8ec
--- /dev/null
+++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerReplicaInfoResult.java
@@ -0,0 +1,44 @@
+/*
+ * 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.hdds.scm.container;
+
+import java.util.List;
+
+/**
+ * Container replica information and the checksum mismatch state determined
+ * by SCM.
+ */
+public final class ContainerReplicaInfoResult {
+
+ private final List replicas;
+ private final boolean dataChecksumMismatch;
+
+ public ContainerReplicaInfoResult(List replicas,
+ boolean dataChecksumMismatch) {
+ this.replicas = replicas;
+ this.dataChecksumMismatch = dataChecksumMismatch;
+ }
+
+ public List getReplicas() {
+ return replicas;
+ }
+
+ public boolean hasDataChecksumMismatch() {
+ return dataChecksumMismatch;
+ }
+}
diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ReplicationManagerReport.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ReplicationManagerReport.java
index aacb5e40dd26..e52caace7520 100644
--- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ReplicationManagerReport.java
+++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/ReplicationManagerReport.java
@@ -90,6 +90,15 @@ public void incrementAndSample(ContainerHealthState stat, ContainerInfo containe
containerHealthState = stat;
}
+ /**
+ * Increments a health state and records the container ID without changing
+ * the primary health state of the container currently being processed.
+ */
+ public void incrementAndSampleAdditionalState(
+ ContainerHealthState stat, ContainerID containerID) {
+ incrementAndSample(stat.name(), containerID);
+ }
+
public void increment(HddsProtos.LifeCycleState stat) {
increment(stat.toString());
}
diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerHealthState.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerHealthState.java
index 6ffc678ea175..f0b5f77d89f2 100644
--- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerHealthState.java
+++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerHealthState.java
@@ -17,6 +17,7 @@
package org.apache.hadoop.hdds.scm.container;
+import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.DATA_CHECKSUM_MISMATCH;
import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.EMPTY;
import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.HEALTHY;
import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.MISSING;
@@ -59,6 +60,7 @@ public void testIndividualStateValues() {
assertEquals(7, OPEN_UNHEALTHY.getValue());
assertEquals(8, QUASI_CLOSED_STUCK.getValue());
assertEquals(9, OPEN_WITHOUT_PIPELINE.getValue());
+ assertEquals(10, DATA_CHECKSUM_MISMATCH.getValue());
}
@Test
@@ -101,6 +103,7 @@ public void testFromValueIndividualStates() {
assertEquals(OPEN_UNHEALTHY, ContainerHealthState.fromValue((short) 7));
assertEquals(QUASI_CLOSED_STUCK, ContainerHealthState.fromValue((short) 8));
assertEquals(OPEN_WITHOUT_PIPELINE, ContainerHealthState.fromValue((short) 9));
+ assertEquals(DATA_CHECKSUM_MISMATCH, ContainerHealthState.fromValue((short) 10));
}
@Test
@@ -136,11 +139,11 @@ public void testAllEnumValuesAreUnique() {
@Test
public void testIndividualStateCount() {
- // Should have 10 individual states (0-9)
+ // Should have 11 individual states (0-10)
long individualCount = java.util.Arrays.stream(ContainerHealthState.values())
.filter(s -> s.getValue() >= 0 && s.getValue() <= 99)
.count();
- assertEquals(10, individualCount, "Expected 10 individual states");
+ assertEquals(11, individualCount, "Expected 11 individual states");
}
@Test
@@ -154,10 +157,10 @@ public void testCombinationStateCount() {
@Test
public void testNoGapsInIndividualValues() {
- // Individual states should be sequential: 0-9
- for (short i = 0; i <= 9; i++) {
+ // Individual states should be sequential: 0-10
+ for (short i = 0; i <= 10; i++) {
ContainerHealthState state = ContainerHealthState.fromValue(i);
- assertTrue(state.getValue() >= 0 && state.getValue() <= 9,
+ assertTrue(state.getValue() >= 0 && state.getValue() <= 10,
"Value " + i + " should map to an individual state");
}
}
diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReplicaChecksumMismatch.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReplicaChecksumMismatch.java
new file mode 100644
index 000000000000..f042c762a5ed
--- /dev/null
+++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReplicaChecksumMismatch.java
@@ -0,0 +1,117 @@
+/*
+ * 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.hdds.scm.container;
+
+import static org.apache.hadoop.hdds.scm.container.ContainerReplicaChecksumMismatch.hasMismatch;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import org.junit.jupiter.api.Test;
+
+class TestContainerReplicaChecksumMismatch {
+
+ @Test
+ void detectsDifferentChecksumsAtTheSameSequenceId() {
+ List replicas = Arrays.asList(
+ new Replica(10L, 100L),
+ new Replica(10L, 200L),
+ new Replica(10L, 100L));
+
+ assertTrue(hasMismatch(replicas, Replica::getSequenceId,
+ Replica::getDataChecksum));
+ }
+
+ @Test
+ void ignoresDifferentChecksumsAtDifferentSequenceIds() {
+ List replicas = Arrays.asList(
+ new Replica(10L, 100L),
+ new Replica(11L, 200L));
+
+ assertFalse(hasMismatch(replicas, Replica::getSequenceId,
+ Replica::getDataChecksum));
+ }
+
+ @Test
+ void ignoresEqualChecksumsAtDifferentSequenceIds() {
+ List replicas = Arrays.asList(
+ new Replica(10L, 100L),
+ new Replica(11L, 100L));
+
+ assertFalse(hasMismatch(replicas, Replica::getSequenceId,
+ Replica::getDataChecksum));
+ }
+
+ @Test
+ void detectsMismatchWithinOneSequenceIdGroup() {
+ List replicas = Arrays.asList(
+ new Replica(10L, 100L),
+ new Replica(11L, 200L),
+ new Replica(11L, 300L));
+
+ assertTrue(hasMismatch(replicas, Replica::getSequenceId,
+ Replica::getDataChecksum));
+ }
+
+ @Test
+ void waitsUntilEveryReplicaReportsADataChecksum() {
+ List replicas = Arrays.asList(
+ new Replica(10L, 100L),
+ new Replica(10L, 200L),
+ new Replica(10L, 0L));
+
+ assertFalse(hasMismatch(replicas, Replica::getSequenceId,
+ Replica::getDataChecksum));
+ }
+
+ @Test
+ void waitsUntilEveryReplicaReportsASequenceId() {
+ List replicas = Arrays.asList(
+ new Replica(10L, 100L),
+ new Replica(null, 200L));
+
+ assertFalse(hasMismatch(replicas, Replica::getSequenceId,
+ Replica::getDataChecksum));
+ }
+
+ @Test
+ void requiresAtLeastTwoReplicas() {
+ assertFalse(hasMismatch(Collections.singletonList(new Replica(10L, 100L)),
+ Replica::getSequenceId, Replica::getDataChecksum));
+ }
+
+ private static final class Replica {
+ private final Long sequenceId;
+ private final long dataChecksum;
+
+ private Replica(Long sequenceId, long dataChecksum) {
+ this.sequenceId = sequenceId;
+ this.dataChecksum = dataChecksum;
+ }
+
+ private Long getSequenceId() {
+ return sequenceId;
+ }
+
+ private long getDataChecksum() {
+ return dataChecksum;
+ }
+ }
+}
diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestReplicationManagerReport.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestReplicationManagerReport.java
index 864c96bae738..9fe3fa931fb0 100644
--- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestReplicationManagerReport.java
+++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/container/TestReplicationManagerReport.java
@@ -28,6 +28,7 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
import java.util.Random;
import java.util.concurrent.ThreadLocalRandom;
@@ -115,6 +116,7 @@ void testJsonOutput() throws IOException {
assertEquals(0, stats.get("OPEN_UNHEALTHY").longValue());
assertEquals(0, stats.get("QUASI_CLOSED_STUCK").longValue());
assertEquals(0, stats.get("OPEN_WITHOUT_PIPELINE").longValue());
+ assertEquals(0, stats.get("DATA_CHECKSUM_MISMATCH").longValue());
JsonNode samples = json.get("samples");
assertEquals(ARRAY, samples.get("UNDER_REPLICATED").getNodeType());
@@ -137,6 +139,20 @@ void testContainerIDsCanBeSampled() {
report.getStat(ContainerHealthState.MIS_REPLICATED));
}
+ @Test
+ void testSampleByContainerIDDoesNotChangeContainerHealthState() {
+ ContainerID containerID = ContainerID.valueOf(1);
+ report.incrementAndSampleAdditionalState(
+ ContainerHealthState.DATA_CHECKSUM_MISMATCH, containerID);
+
+ assertEquals(1,
+ report.getStat(ContainerHealthState.DATA_CHECKSUM_MISMATCH));
+ assertEquals(Collections.singletonList(containerID),
+ report.getSample(ContainerHealthState.DATA_CHECKSUM_MISMATCH));
+ assertEquals(ContainerHealthState.HEALTHY,
+ report.getContainerHealthState());
+ }
+
@Test
void testSamplesAreLimited() {
verifySampleLimit(report, 100);
diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
index 9a41a047ee1a..85f3179a72d6 100644
--- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
+++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
@@ -38,6 +38,7 @@
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
import org.apache.hadoop.hdds.scm.container.ContainerListResult;
import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfo;
+import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfoResult;
import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
import org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
@@ -83,6 +84,17 @@ ContainerWithPipeline getContainerWithPipeline(long containerId)
List getContainerReplicas(
long containerId) throws IOException;
+ /**
+ * Gets replica information and the container-level status determined by
+ * SCM.
+ *
+ * @param containerId the container ID
+ * @return replicas and container-level status
+ * @throws IOException on failure
+ */
+ ContainerReplicaInfoResult getContainerReplicasWithStatus(long containerId)
+ throws IOException;
+
/**
* Close a container.
*
diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java
index 70b758ef1c44..f7df1653f0ca 100644
--- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java
+++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java
@@ -34,6 +34,7 @@
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DeletedBlocksTransactionSummary;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.DecommissionScmResponseProto;
+import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerReplicasResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.SafeModeRuleStatusProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.StartContainerBalancerResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.Type;
@@ -113,6 +114,17 @@ ContainerWithPipeline getContainerWithPipeline(long containerID)
List getContainerReplicas(
long containerId, int clientVersion) throws IOException;
+ /**
+ * Gets the replicas and container-level status returned by SCM.
+ *
+ * @param containerId ID of the container
+ * @param clientVersion client version
+ * @return response containing replicas and container-level status
+ * @throws IOException on failure
+ */
+ GetContainerReplicasResponseProto getContainerReplicasResponse(
+ long containerId, int clientVersion) throws IOException;
+
/**
* Ask SCM the location of a batch of containers. SCM responds with a group of
* nodes where these containers and their replicas are located.
diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java
index 7cdb5fa33652..7439af086fc2 100644
--- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java
+++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java
@@ -74,6 +74,7 @@
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerCountRequestProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerCountResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerReplicasRequestProto;
+import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerReplicasResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerRequestProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerTokenRequestProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerTokenResponseProto;
@@ -340,6 +341,13 @@ public ContainerWithPipeline getContainerWithPipeline(long containerID)
@Override
public List getContainerReplicas(
long containerID, int clientVersion) throws IOException {
+ return getContainerReplicasResponse(containerID, clientVersion)
+ .getContainerReplicaList();
+ }
+
+ @Override
+ public GetContainerReplicasResponseProto getContainerReplicasResponse(
+ long containerID, int clientVersion) throws IOException {
Preconditions.checkState(containerID >= 0,
"Container ID cannot be negative");
@@ -351,7 +359,7 @@ public List getContainerReplicas(
ScmContainerLocationResponse response =
submitRequest(Type.GetContainerReplicas,
(builder) -> builder.setGetContainerReplicasRequest(request));
- return response.getGetContainerReplicasResponse().getContainerReplicaList();
+ return response.getGetContainerReplicasResponse();
}
/**
diff --git a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
index 9e7244ecf828..30a435d9ec07 100644
--- a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
+++ b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
@@ -264,6 +264,7 @@ message GetContainerReplicasRequestProto {
message GetContainerReplicasResponseProto {
repeated SCMContainerReplicaProto containerReplica = 1;
+ optional bool dataChecksumMismatch = 2;
}
message GetContainerWithPipelineBatchRequestProto {
diff --git a/hadoop-hdds/interface-admin/src/main/resources/proto.lock b/hadoop-hdds/interface-admin/src/main/resources/proto.lock
index 749549070093..51a72a5e7f92 100644
--- a/hadoop-hdds/interface-admin/src/main/resources/proto.lock
+++ b/hadoop-hdds/interface-admin/src/main/resources/proto.lock
@@ -1066,6 +1066,12 @@
"name": "containerReplica",
"type": "SCMContainerReplicaProto",
"is_repeated": true
+ },
+ {
+ "id": 2,
+ "name": "dataChecksumMismatch",
+ "type": "bool",
+ "optional": true
}
]
},
@@ -2590,4 +2596,4 @@
}
}
]
-}
\ No newline at end of file
+}
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
index f890fb6a082a..919e857f7128 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
@@ -63,6 +63,7 @@
import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
import org.apache.hadoop.hdds.scm.container.replication.health.ClosedWithUnhealthyReplicasHandler;
import org.apache.hadoop.hdds.scm.container.replication.health.ClosingContainerHandler;
+import org.apache.hadoop.hdds.scm.container.replication.health.DataChecksumMismatchCheckHandler;
import org.apache.hadoop.hdds.scm.container.replication.health.DeletingContainerHandler;
import org.apache.hadoop.hdds.scm.container.replication.health.ECMisReplicationCheckHandler;
import org.apache.hadoop.hdds.scm.container.replication.health.ECReplicationCheckHandler;
@@ -169,6 +170,7 @@ public class ReplicationManager implements SCMService, ContainerReplicaPendingOp
private long lastTimeToBeReadyInMillis = 0;
private final Clock clock;
private final ContainerReplicaPendingOps containerReplicaPendingOps;
+ private final DataChecksumMismatchCheckHandler checksumMismatchCheckHandler;
private final ECReplicationCheckHandler ecReplicationCheckHandler;
private final ECMisReplicationCheckHandler ecMisReplicationCheckHandler;
private final RatisReplicationCheckHandler ratisReplicationCheckHandler;
@@ -231,6 +233,8 @@ public ReplicationManager(final ReplicationManagerConfiguration rmConf,
HddsConfigKeys.HDDS_SCM_WAIT_TIME_AFTER_SAFE_MODE_EXIT_DEFAULT,
TimeUnit.MILLISECONDS);
this.containerReplicaPendingOps = replicaPendingOps;
+ this.checksumMismatchCheckHandler =
+ new DataChecksumMismatchCheckHandler();
this.ecReplicationCheckHandler = new ECReplicationCheckHandler();
this.ecMisReplicationCheckHandler =
new ECMisReplicationCheckHandler(ecContainerPlacement);
@@ -270,6 +274,7 @@ public ReplicationManager(final ReplicationManagerConfiguration rmConf,
.addNext(new DeletingContainerHandler(this))
.addNext(new QuasiClosedStuckReplicationCheck(rmConf))
.addNext(ecReplicationCheckHandler)
+ .addNext(checksumMismatchCheckHandler)
.addNext(ratisReplicationCheckHandler)
.addNext(new ClosedWithUnhealthyReplicasHandler(this))
.addNext(ecMisReplicationCheckHandler)
@@ -367,6 +372,8 @@ public synchronized void processAll() {
ReplicationManagerReport report = new ReplicationManagerReport(
rmConf.getContainerSampleLimit());
ReplicationQueue newRepQueue = new ReplicationQueue();
+ checksumMismatchCheckHandler.startScan();
+ int processedContainers = 0;
for (ContainerInfo c : containers) {
if (!shouldRun()) {
break;
@@ -378,6 +385,12 @@ public synchronized void processAll() {
} catch (ContainerNotFoundException e) {
LOG.error("Container {} not found", c.getContainerID(), e);
}
+ processedContainers++;
+ }
+ if (processedContainers == containers.size()) {
+ checksumMismatchCheckHandler.completeScan(report);
+ } else {
+ checksumMismatchCheckHandler.abortScan();
}
report.setComplete();
replicationQueue.set(newRepQueue);
@@ -1482,6 +1495,10 @@ public ReplicationManagerMetrics getMetrics() {
return metrics;
}
+ public boolean hasContainerChecksumMismatch(ContainerID containerID) {
+ return checksumMismatchCheckHandler.hasPersistentMismatch(containerID);
+ }
+
public ReplicationManagerConfiguration getConfig() {
return rmConf;
}
@@ -1591,4 +1608,3 @@ public synchronized boolean notifyNodeStateChange() {
}
}
}
-
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/health/DataChecksumMismatchCheckHandler.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/health/DataChecksumMismatchCheckHandler.java
new file mode 100644
index 000000000000..713da291d50b
--- /dev/null
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/health/DataChecksumMismatchCheckHandler.java
@@ -0,0 +1,131 @@
+/*
+ * 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.hdds.scm.container.replication.health;
+
+import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState.CLOSED;
+import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType.RATIS;
+import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.DATA_CHECKSUM_MISMATCH;
+import static org.apache.hadoop.hdds.scm.container.ContainerReplicaChecksumMismatch.hasMismatch;
+
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.HddsUtils;
+import org.apache.hadoop.hdds.scm.container.ContainerID;
+import org.apache.hadoop.hdds.scm.container.ContainerInfo;
+import org.apache.hadoop.hdds.scm.container.ContainerReplica;
+import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
+import org.apache.hadoop.hdds.scm.container.replication.ContainerCheckRequest;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Detects replicas with the same BCSID and different data checksums.
+ */
+public class DataChecksumMismatchCheckHandler extends AbstractCheck {
+
+ private static final Logger LOG =
+ LoggerFactory.getLogger(DataChecksumMismatchCheckHandler.class);
+
+ private final Map> currentMismatches =
+ new HashMap<>();
+ private final Set mismatchesInPreviousScan = new HashSet<>();
+ private final Set warnedMismatches = new HashSet<>();
+ private volatile Set persistentMismatches =
+ Collections.emptySet();
+ private boolean scanInProgress;
+
+ @Override
+ public boolean handle(ContainerCheckRequest request) {
+ ContainerInfo container = request.getContainerInfo();
+ Set replicas = request.getContainerReplicas();
+ if (container.getState() == CLOSED &&
+ container.getReplicationType() == RATIS &&
+ hasMismatch(replicas, ContainerReplica::getSequenceId,
+ ContainerReplica::getDataChecksum)) {
+ if (request.isReadOnly()) {
+ if (hasPersistentMismatch(container.containerID())) {
+ request.getReport().incrementAndSampleAdditionalState(
+ DATA_CHECKSUM_MISMATCH, container.containerID());
+ }
+ } else if (scanInProgress) {
+ currentMismatches.put(container.containerID(), replicas);
+ }
+ }
+ return false;
+ }
+
+ /** Starts collecting mismatches for a new Replication Manager scan. */
+ public void startScan() {
+ currentMismatches.clear();
+ scanInProgress = true;
+ }
+
+ /**
+ * Completes a scan and reports mismatches seen in two consecutive scans.
+ */
+ public void completeScan(ReplicationManagerReport report) {
+ Set confirmed = new HashSet<>(currentMismatches.keySet());
+ confirmed.retainAll(mismatchesInPreviousScan);
+
+ confirmed.stream()
+ .sorted(Comparator.comparingLong(ContainerID::getId))
+ .forEach(containerID -> report.incrementAndSampleAdditionalState(
+ DATA_CHECKSUM_MISMATCH, containerID));
+
+ Set newWarnings = new HashSet<>(confirmed);
+ newWarnings.removeAll(warnedMismatches);
+ newWarnings.stream()
+ .sorted(Comparator.comparingLong(ContainerID::getId))
+ .forEach(containerID -> LOG.warn(
+ "Container {} has replicas with the same BCSID but different data checksums: {}",
+ containerID, formatChecksumDetails(currentMismatches.get(containerID))));
+
+ warnedMismatches.retainAll(currentMismatches.keySet());
+ warnedMismatches.addAll(confirmed);
+ mismatchesInPreviousScan.clear();
+ mismatchesInPreviousScan.addAll(currentMismatches.keySet());
+ persistentMismatches = Collections.unmodifiableSet(confirmed);
+ currentMismatches.clear();
+ scanInProgress = false;
+ }
+
+ /** Discards observations when a scan does not process every container. */
+ public void abortScan() {
+ currentMismatches.clear();
+ scanInProgress = false;
+ }
+
+ public boolean hasPersistentMismatch(ContainerID containerID) {
+ return persistentMismatches.contains(containerID);
+ }
+
+ private static String formatChecksumDetails(Set replicas) {
+ return replicas.stream()
+ .sorted()
+ .map(replica -> replica.getDatanodeDetails().getUuidString() +
+ "(BCSID=" + replica.getSequenceId() +
+ ", dataChecksum=" +
+ HddsUtils.checksumToString(replica.getDataChecksum()) + ")")
+ .collect(Collectors.joining(", "));
+ }
+}
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
index c1561f0cd19c..e4f32a2285a1 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
@@ -781,10 +781,8 @@ public ScmContainerLocationResponse processRequest(
public GetContainerReplicasResponseProto getContainerReplicas(
GetContainerReplicasRequestProto request, int clientVersion)
throws IOException {
- List replicas
- = impl.getContainerReplicas(request.getContainerID(), clientVersion);
- return GetContainerReplicasResponseProto.newBuilder()
- .addAllContainerReplica(replicas).build();
+ return impl.getContainerReplicasResponse(request.getContainerID(),
+ clientVersion);
}
public ContainerResponseProto allocateContainer(ContainerRequestProto request,
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
index 8069971459ec..240af7713ff8 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java
@@ -64,6 +64,7 @@
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.DecommissionScmResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.DecommissionScmResponseProto.Builder;
+import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerReplicasResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.SafeModeRuleStatusProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.StartContainerBalancerResponseProto;
import org.apache.hadoop.hdds.protocolPB.ReconfigureProtocolPB;
@@ -387,6 +388,16 @@ public List getContainerReplicas(
}
}
+ @Override
+ public GetContainerReplicasResponseProto getContainerReplicasResponse(
+ long containerId, int clientVersion) throws IOException {
+ return GetContainerReplicasResponseProto.newBuilder()
+ .addAllContainerReplica(getContainerReplicas(containerId, clientVersion))
+ .setDataChecksumMismatch(getScm().getReplicationManager()
+ .hasContainerChecksumMismatch(ContainerID.valueOf(containerId)))
+ .build();
+ }
+
@Override
public List getContainerWithPipelineBatch(
Iterable extends Long> containerIDs) throws IOException {
diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
index bcb3dea97680..f83d1fe21adb 100644
--- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
+++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
@@ -23,6 +23,7 @@
import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE;
import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
import static org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicationCommandPriority.LOW;
+import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.DATA_CHECKSUM_MISMATCH;
import static org.apache.hadoop.hdds.scm.container.replication.ContainerReplicaOp.PendingOpType.ADD;
import static org.apache.hadoop.hdds.scm.container.replication.ReplicationTestUtil.createContainerInfo;
import static org.apache.hadoop.hdds.scm.container.replication.ReplicationTestUtil.createContainerReplica;
@@ -64,6 +65,7 @@
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
+import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.hadoop.hdds.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
@@ -77,6 +79,7 @@
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
import org.apache.hadoop.hdds.scm.HddsTestUtils;
import org.apache.hadoop.hdds.scm.PlacementPolicy;
+import org.apache.hadoop.hdds.scm.container.ContainerChecksums;
import org.apache.hadoop.hdds.scm.container.ContainerHealthState;
import org.apache.hadoop.hdds.scm.container.ContainerID;
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
@@ -85,6 +88,7 @@
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
import org.apache.hadoop.hdds.scm.container.placement.algorithms.ContainerPlacementStatusDefault;
+import org.apache.hadoop.hdds.scm.container.replication.health.DataChecksumMismatchCheckHandler;
import org.apache.hadoop.hdds.scm.events.SCMEvents;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.ha.SCMContext;
@@ -101,6 +105,7 @@
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.SCMCommand;
import org.apache.ozone.test.GenericTestUtils;
+import org.apache.ozone.test.GenericTestUtils.LogCapturer;
import org.apache.ozone.test.MockClock;
import org.apache.ratis.protocol.exceptions.NotLeaderException;
import org.junit.jupiter.api.AfterEach;
@@ -1084,6 +1089,66 @@ public void testUnderReplicationQueuePopulated() {
assertNull(res);
}
+ @Test
+ public void testDataChecksumMismatchIsDebouncedAndWarnedOnce()
+ throws ContainerNotFoundException {
+ RatisReplicationConfig ratisReplicationConfig =
+ RatisReplicationConfig.getInstance(THREE);
+ ContainerInfo container = createContainerInfo(ratisReplicationConfig, 10,
+ HddsProtos.LifeCycleState.CLOSED);
+ Set mismatch = createReplicasWithChecksums(container,
+ new long[]{10, 10, 10}, new long[]{100, 200, 100});
+ containerInfoSet.add(container);
+ containerReplicaMap.put(container.containerID(), mismatch);
+ enableProcessAll();
+
+ LogCapturer logs =
+ LogCapturer.captureLogs(DataChecksumMismatchCheckHandler.class);
+ String warning = "same BCSID but different data checksums";
+ clearInvocations(containerManager);
+ replicationManager.processAll();
+ verify(containerManager, times(1))
+ .getContainerReplicas(container.containerID());
+ assertEquals(0, replicationManager.getContainerReport()
+ .getStat(DATA_CHECKSUM_MISMATCH));
+ assertFalse(replicationManager.hasContainerChecksumMismatch(
+ container.containerID()));
+ assertEquals(0, StringUtils.countMatches(logs.getOutput(), warning));
+
+ replicationManager.processAll();
+ assertEquals(1, replicationManager.getContainerReport()
+ .getStat(DATA_CHECKSUM_MISMATCH));
+ assertThat(replicationManager.getContainerReport()
+ .getSample(DATA_CHECKSUM_MISMATCH))
+ .containsExactly(container.containerID());
+ assertTrue(replicationManager.hasContainerChecksumMismatch(
+ container.containerID()));
+ assertEquals(1, StringUtils.countMatches(logs.getOutput(), warning));
+
+ replicationManager.processAll();
+ assertEquals(1, replicationManager.getContainerReport()
+ .getStat(DATA_CHECKSUM_MISMATCH));
+ assertEquals(1, StringUtils.countMatches(logs.getOutput(), warning));
+
+ containerReplicaMap.put(container.containerID(),
+ createReplicasWithChecksums(container, new long[]{10, 10, 10},
+ new long[]{100, 100, 100}));
+ replicationManager.processAll();
+ assertEquals(0, replicationManager.getContainerReport()
+ .getStat(DATA_CHECKSUM_MISMATCH));
+ assertFalse(replicationManager.hasContainerChecksumMismatch(
+ container.containerID()));
+
+ containerReplicaMap.put(container.containerID(), mismatch);
+ replicationManager.processAll();
+ assertEquals(0, replicationManager.getContainerReport()
+ .getStat(DATA_CHECKSUM_MISMATCH));
+ replicationManager.processAll();
+ assertEquals(1, replicationManager.getContainerReport()
+ .getStat(DATA_CHECKSUM_MISMATCH));
+ assertEquals(2, StringUtils.countMatches(logs.getOutput(), warning));
+ }
+
@Test
public void testSendDatanodeDeleteCommand() throws NotLeaderException {
ECReplicationConfig ecRepConfig = new ECReplicationConfig(3, 2);
@@ -1697,6 +1762,21 @@ public void testReconfigureContainerSampleLimit() {
"Second report should have 50 samples after reconfiguration");
}
+ private static Set createReplicasWithChecksums(
+ ContainerInfo container, long[] sequenceIds, long[] dataChecksums) {
+ assertEquals(sequenceIds.length, dataChecksums.length);
+ Set replicas = new HashSet<>();
+ for (int i = 0; i < sequenceIds.length; i++) {
+ ContainerReplica replica = createContainerReplica(container.containerID(),
+ 0, IN_SERVICE, ContainerReplicaProto.State.CLOSED,
+ sequenceIds[i]).toBuilder()
+ .setChecksums(ContainerChecksums.of(dataChecksums[i]))
+ .build();
+ replicas.add(replica);
+ }
+ return replicas;
+ }
+
@SafeVarargs
private final Set addReplicas(ContainerInfo container,
ContainerReplicaProto.State replicaState,
diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/health/TestDataChecksumMismatchCheckHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/health/TestDataChecksumMismatchCheckHandler.java
new file mode 100644
index 000000000000..03f65742fa7e
--- /dev/null
+++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/health/TestDataChecksumMismatchCheckHandler.java
@@ -0,0 +1,181 @@
+/*
+ * 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.hdds.scm.container.replication.health;
+
+import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE;
+import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
+import static org.apache.hadoop.hdds.scm.container.ContainerHealthState.DATA_CHECKSUM_MISMATCH;
+import static org.apache.hadoop.hdds.scm.container.replication.ReplicationTestUtil.createContainerInfo;
+import static org.apache.hadoop.hdds.scm.container.replication.ReplicationTestUtil.createContainerReplica;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.Set;
+import org.apache.hadoop.hdds.client.ECReplicationConfig;
+import org.apache.hadoop.hdds.client.RatisReplicationConfig;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
+import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto;
+import org.apache.hadoop.hdds.scm.container.ContainerChecksums;
+import org.apache.hadoop.hdds.scm.container.ContainerInfo;
+import org.apache.hadoop.hdds.scm.container.ContainerReplica;
+import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
+import org.apache.hadoop.hdds.scm.container.replication.ContainerCheckRequest;
+import org.apache.hadoop.hdds.scm.container.replication.ReplicationQueue;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Tests for {@link DataChecksumMismatchCheckHandler}.
+ */
+public class TestDataChecksumMismatchCheckHandler {
+
+ private DataChecksumMismatchCheckHandler handler;
+
+ @BeforeEach
+ void setup() {
+ handler = new DataChecksumMismatchCheckHandler();
+ }
+
+ @Test
+ void testMismatchIsReportedAfterTwoScans() {
+ ContainerInfo container = createContainerInfo(
+ RatisReplicationConfig.getInstance(THREE), 10,
+ HddsProtos.LifeCycleState.CLOSED);
+ Set replicas = createReplicasWithChecksums(container,
+ new long[]{10, 10, 10}, new long[]{100, 200, 100});
+
+ ReplicationManagerReport firstReport = runScan(container, replicas);
+ assertEquals(0, firstReport.getStat(DATA_CHECKSUM_MISMATCH));
+ assertFalse(handler.hasPersistentMismatch(container.containerID()));
+
+ ReplicationManagerReport secondReport = runScan(container, replicas);
+ assertEquals(1, secondReport.getStat(DATA_CHECKSUM_MISMATCH));
+ assertThat(secondReport.getSample(DATA_CHECKSUM_MISMATCH))
+ .containsExactly(container.containerID());
+ assertThat(handler.hasPersistentMismatch(container.containerID())).isTrue();
+
+ ReplicationManagerReport resolvedReport = runScan(container,
+ createReplicasWithChecksums(container, new long[]{10, 10, 10},
+ new long[]{100, 100, 100}));
+ assertEquals(0, resolvedReport.getStat(DATA_CHECKSUM_MISMATCH));
+ assertFalse(handler.hasPersistentMismatch(container.containerID()));
+ }
+
+ @Test
+ void testOnlyClosedRatisContainersAreChecked() {
+ ContainerInfo openRatis = createContainerInfo(
+ RatisReplicationConfig.getInstance(THREE), 10,
+ HddsProtos.LifeCycleState.OPEN);
+ Set replicas = createReplicasWithChecksums(openRatis,
+ new long[]{10, 10, 10}, new long[]{100, 200, 100});
+ runScan(openRatis, replicas);
+ ReplicationManagerReport report = runScan(openRatis, replicas);
+ assertEquals(0, report.getStat(DATA_CHECKSUM_MISMATCH));
+
+ ContainerInfo closedEc = createContainerInfo(
+ new ECReplicationConfig(3, 2), 10,
+ HddsProtos.LifeCycleState.CLOSED);
+ replicas = createReplicasWithChecksums(closedEc,
+ new long[]{10, 10, 10}, new long[]{100, 200, 100});
+ runScan(closedEc, replicas);
+ report = runScan(closedEc, replicas);
+ assertEquals(0, report.getStat(DATA_CHECKSUM_MISMATCH));
+ }
+
+ @Test
+ void testIncompleteScanIsIgnored() {
+ ContainerInfo container = createContainerInfo(
+ RatisReplicationConfig.getInstance(THREE), 10,
+ HddsProtos.LifeCycleState.CLOSED);
+ Set replicas = createReplicasWithChecksums(container,
+ new long[]{10, 10, 10}, new long[]{100, 200, 100});
+
+ handler.startScan();
+ assertFalse(handler.handle(createRequest(container, replicas,
+ new ReplicationManagerReport(10))));
+ handler.abortScan();
+
+ ReplicationManagerReport report = runScan(container, replicas);
+ assertEquals(0, report.getStat(DATA_CHECKSUM_MISMATCH));
+ assertFalse(handler.hasPersistentMismatch(container.containerID()));
+ }
+
+ @Test
+ void testReadOnlyCheckDoesNotAffectDebounce() {
+ ContainerInfo container = createContainerInfo(
+ RatisReplicationConfig.getInstance(THREE), 10,
+ HddsProtos.LifeCycleState.CLOSED);
+ Set replicas = createReplicasWithChecksums(container,
+ new long[]{10, 10, 10}, new long[]{100, 200, 100});
+
+ handler.startScan();
+ ReplicationManagerReport readOnlyReport = new ReplicationManagerReport(10);
+ assertFalse(handler.handle(createRequest(container, replicas,
+ readOnlyReport, true)));
+ handler.completeScan(readOnlyReport);
+
+ ReplicationManagerReport report = runScan(container, replicas);
+ assertEquals(0, report.getStat(DATA_CHECKSUM_MISMATCH));
+ assertFalse(handler.hasPersistentMismatch(container.containerID()));
+ }
+
+ private ReplicationManagerReport runScan(ContainerInfo container,
+ Set replicas) {
+ ReplicationManagerReport report = new ReplicationManagerReport(10);
+ handler.startScan();
+ assertFalse(handler.handle(createRequest(container, replicas, report)));
+ handler.completeScan(report);
+ return report;
+ }
+
+ private static ContainerCheckRequest createRequest(ContainerInfo container,
+ Set replicas, ReplicationManagerReport report) {
+ return createRequest(container, replicas, report, false);
+ }
+
+ private static ContainerCheckRequest createRequest(ContainerInfo container,
+ Set replicas, ReplicationManagerReport report,
+ boolean readOnly) {
+ return new ContainerCheckRequest.Builder()
+ .setContainerInfo(container)
+ .setContainerReplicas(replicas)
+ .setPendingOps(Collections.emptyList())
+ .setMaintenanceRedundancy(2)
+ .setReport(report)
+ .setReplicationQueue(new ReplicationQueue())
+ .setReadOnly(readOnly)
+ .build();
+ }
+
+ private static Set createReplicasWithChecksums(
+ ContainerInfo container, long[] sequenceIds, long[] dataChecksums) {
+ assertEquals(sequenceIds.length, dataChecksums.length);
+ Set replicas = new HashSet<>();
+ for (int i = 0; i < sequenceIds.length; i++) {
+ replicas.add(createContainerReplica(container.containerID(), i + 1,
+ IN_SERVICE, ContainerReplicaProto.State.CLOSED, sequenceIds[i])
+ .toBuilder()
+ .setChecksums(ContainerChecksums.of(dataChecksums[i]))
+ .build());
+ }
+ return replicas;
+ }
+}
diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java
index 7fee21620d15..8d6b9d874976 100644
--- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java
+++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java
@@ -36,6 +36,7 @@
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DeletedBlocksTransactionSummary;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.DecommissionScmResponseProto;
+import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerReplicasResponseProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.SafeModeRuleStatusProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.StartContainerBalancerResponseProto;
import org.apache.hadoop.hdds.scm.DatanodeAdminError;
@@ -48,6 +49,7 @@
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
import org.apache.hadoop.hdds.scm.container.ContainerListResult;
import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfo;
+import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfoResult;
import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
import org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
@@ -350,14 +352,21 @@ public ContainerWithPipeline getContainerWithPipeline(long containerId)
@Override
public List getContainerReplicas(long containerId) throws IOException {
- List protos =
- storageContainerLocationClient.getContainerReplicas(containerId,
- ClientVersion.CURRENT_VERSION);
+ return getContainerReplicasWithStatus(containerId).getReplicas();
+ }
+
+ @Override
+ public ContainerReplicaInfoResult getContainerReplicasWithStatus(
+ long containerId) throws IOException {
+ GetContainerReplicasResponseProto response = storageContainerLocationClient
+ .getContainerReplicasResponse(containerId, ClientVersion.CURRENT_VERSION);
List replicas = new ArrayList<>();
- for (HddsProtos.SCMContainerReplicaProto p : protos) {
+ for (HddsProtos.SCMContainerReplicaProto p :
+ response.getContainerReplicaList()) {
replicas.add(ContainerReplicaInfo.fromProto(p));
}
- return replicas;
+ return new ContainerReplicaInfoResult(replicas,
+ response.getDataChecksumMismatch());
}
@Override
diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/container/InfoSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/container/InfoSubcommand.java
index ea296ba98703..9bf6367b8467 100644
--- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/container/InfoSubcommand.java
+++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/container/InfoSubcommand.java
@@ -35,6 +35,7 @@
import org.apache.hadoop.hdds.scm.client.ScmClient;
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfo;
+import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfoResult;
import org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
import org.apache.hadoop.hdds.scm.ha.SCMHAUtils;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
@@ -118,8 +119,12 @@ private void printDetails(ScmClient scmClient, long containerID) throws IOExcept
}
List replicas = null;
+ boolean dataChecksumMismatch = false;
try {
- replicas = scmClient.getContainerReplicas(containerID);
+ ContainerReplicaInfoResult result =
+ scmClient.getContainerReplicasWithStatus(containerID);
+ replicas = result.getReplicas();
+ dataChecksumMismatch = result.hasDataChecksumMismatch();
} catch (IOException e) {
rootCommand().printError(e);
}
@@ -129,13 +134,15 @@ private void printDetails(ScmClient scmClient, long containerID) throws IOExcept
ContainerWithPipelineAndReplicas wrapper =
new ContainerWithPipelineAndReplicas(container.getContainerInfo(),
container.getPipeline(), replicas,
- container.getContainerInfo().getPipelineID());
+ container.getContainerInfo().getPipelineID(),
+ dataChecksumMismatch);
System.out.println(JsonUtils.toJsonStringWithDefaultPrettyPrinter(wrapper));
} else {
ContainerWithoutDatanodes wrapper =
new ContainerWithoutDatanodes(container.getContainerInfo(),
container.getPipeline(), replicas,
- container.getContainerInfo().getPipelineID());
+ container.getContainerInfo().getPipelineID(),
+ dataChecksumMismatch);
System.out.println(JsonUtils.toJsonStringWithDefaultPrettyPrinter(wrapper));
}
} else {
@@ -175,6 +182,9 @@ private void printDetails(ScmClient scmClient, long containerID) throws IOExcept
// Print the replica details if available
if (replicas != null) {
+ if (dataChecksumMismatch) {
+ System.out.println("Data checksum mismatch: true");
+ }
String replicaStr = replicas.stream()
.sorted(Comparator.comparing(ContainerReplicaInfo::getReplicaIndex))
.map(InfoSubcommand::buildReplicaDetails)
@@ -195,6 +205,8 @@ private static String buildReplicaDetails(ContainerReplicaInfo replica) {
sb.append(" ReplicaIndex: ").append(replica.getReplicaIndex()).append(';');
}
sb.append(" SequenceId: ").append(replica.getSequenceId()).append(';')
+ .append(" DataChecksum: ")
+ .append(HddsUtils.checksumToString(replica.getDataChecksum())).append(';')
.append(" Origin: ").append(replica.getPlaceOfBirth().toString()).append(';')
.append(" Location: ").append(buildDatanodeDetails(replica.getDatanodeDetails()));
return sb.toString();
@@ -206,13 +218,16 @@ private static class ContainerWithPipelineAndReplicas {
private Pipeline pipeline;
private List replicas;
private PipelineID writePipelineID;
+ private boolean dataChecksumMismatch;
ContainerWithPipelineAndReplicas(ContainerInfo container, Pipeline pipeline,
- List replicas, PipelineID pipelineID) {
+ List replicas, PipelineID pipelineID,
+ boolean dataChecksumMismatch) {
this.containerInfo = container;
this.pipeline = pipeline;
this.replicas = replicas;
this.writePipelineID = pipelineID;
+ this.dataChecksumMismatch = dataChecksumMismatch;
}
public ContainerInfo getContainerInfo() {
@@ -231,6 +246,10 @@ public PipelineID getWritePipelineID() {
return writePipelineID;
}
+ public boolean getDataChecksumMismatch() {
+ return dataChecksumMismatch;
+ }
+
}
private static class ContainerWithoutDatanodes {
@@ -239,13 +258,16 @@ private static class ContainerWithoutDatanodes {
private PipelineWithoutDatanodes pipeline;
private List replicas;
private PipelineID writePipelineId;
+ private boolean dataChecksumMismatch;
ContainerWithoutDatanodes(ContainerInfo container, Pipeline pipeline,
- List replicas, PipelineID pipelineID) {
+ List replicas, PipelineID pipelineID,
+ boolean dataChecksumMismatch) {
this.containerInfo = container;
this.pipeline = new PipelineWithoutDatanodes(pipeline);
this.replicas = replicas;
this.writePipelineId = pipelineID;
+ this.dataChecksumMismatch = dataChecksumMismatch;
}
public ContainerInfo getContainerInfo() {
@@ -263,6 +285,10 @@ public List getReplicas() {
public PipelineID getWritePipelineId() {
return writePipelineId;
}
+
+ public boolean getDataChecksumMismatch() {
+ return dataChecksumMismatch;
+ }
}
// All Pipeline information except the ones dependent on datanodes
diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/container/TestInfoSubCommand.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/container/TestInfoSubCommand.java
index 86e3bfdf2731..5567c14e22ea 100644
--- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/container/TestInfoSubCommand.java
+++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/container/TestInfoSubCommand.java
@@ -26,6 +26,8 @@
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.any;
import static org.mockito.Mockito.anyLong;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -43,6 +45,7 @@
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.HddsUtils;
import org.apache.hadoop.hdds.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
@@ -50,6 +53,7 @@
import org.apache.hadoop.hdds.scm.client.ScmClient;
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfo;
+import org.apache.hadoop.hdds.scm.container.ContainerReplicaInfoResult;
import org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
@@ -81,6 +85,9 @@ public void setup() throws IOException {
scmClient = mock(ScmClient.class);
datanodes = createDatanodeDetails(3);
when(scmClient.getContainerWithPipeline(anyLong())).then(i -> getContainerWithPipeline(i.getArgument(0)));
+ doAnswer(invocation -> new ContainerReplicaInfoResult(
+ scmClient.getContainerReplicas(invocation.getArgument(0)), false))
+ .when(scmClient).getContainerReplicasWithStatus(anyLong());
when(scmClient.getPipeline(any())).thenThrow(new PipelineNotFoundException("Pipeline not found."));
System.setOut(new PrintStream(outContent, false, DEFAULT_ENCODING));
@@ -104,6 +111,50 @@ public void testReplicaIndexInOutput() throws Exception {
testReplicaIncludedInOutput(true);
}
+ @Test
+ public void testDataChecksumMismatchInOutput() throws Exception {
+ doReturn(new ContainerReplicaInfoResult(
+ getReplicas(false, new long[]{10, 10, 10},
+ new long[]{100, 200, 100}), true))
+ .when(scmClient).getContainerReplicasWithStatus(anyLong());
+ cmd = new InfoSubcommand();
+ new CommandLine(cmd).parseArgs("1");
+ cmd.execute(scmClient);
+
+ assertThat(outContent.toString(DEFAULT_ENCODING))
+ .contains("Data checksum mismatch: true")
+ .contains("DataChecksum: " + HddsUtils.checksumToString(100));
+ }
+
+ @Test
+ public void testUnconfirmedDataChecksumMismatchIsNotInOutput()
+ throws Exception {
+ doReturn(new ContainerReplicaInfoResult(
+ getReplicas(false, new long[]{10, 10, 10},
+ new long[]{100, 200, 100}), false))
+ .when(scmClient).getContainerReplicasWithStatus(anyLong());
+ cmd = new InfoSubcommand();
+ new CommandLine(cmd).parseArgs("1");
+ cmd.execute(scmClient);
+
+ assertThat(outContent.toString(DEFAULT_ENCODING))
+ .doesNotContain("Data checksum mismatch: true");
+ }
+
+ @Test
+ public void testDataChecksumMismatchInJsonOutput() throws Exception {
+ doReturn(new ContainerReplicaInfoResult(
+ getReplicas(false, new long[]{10, 10, 10},
+ new long[]{100, 200, 100}), true))
+ .when(scmClient).getContainerReplicasWithStatus(anyLong());
+ cmd = new InfoSubcommand();
+ new CommandLine(cmd).parseArgs("1", "--json");
+ cmd.execute(scmClient);
+
+ assertThat(outContent.toString(DEFAULT_ENCODING))
+ .contains("\"dataChecksumMismatch\" : true");
+ }
+
@Test
public void testErrorWhenNoContainerIDParam() throws Exception {
cmd = new InfoSubcommand();
@@ -332,9 +383,18 @@ private void testJsonOutput() throws IOException {
}
private List getReplicas(boolean includeIndex) {
+ return getReplicas(includeIndex, new long[]{1, 1, 1},
+ new long[]{0, 0, 0});
+ }
+
+ private List getReplicas(boolean includeIndex,
+ long[] sequenceIds, long[] dataChecksums) {
+ assertEquals(datanodes.size(), sequenceIds.length);
+ assertEquals(datanodes.size(), dataChecksums.length);
List replicas = new ArrayList<>();
int index = 1;
- for (DatanodeDetails dn : datanodes) {
+ for (int i = 0; i < datanodes.size(); i++) {
+ DatanodeDetails dn = datanodes.get(i);
ContainerReplicaInfo.Builder container
= new ContainerReplicaInfo.Builder()
.setContainerID(1)
@@ -343,7 +403,8 @@ private List getReplicas(boolean includeIndex) {
.setPlaceOfBirth(dn.getID())
.setDatanodeDetails(dn)
.setKeyCount(1)
- .setSequenceId(1);
+ .setSequenceId(sequenceIds[i])
+ .setDataChecksum(dataChecksums[i]);
if (includeIndex) {
container.setReplicaIndex(index++);
}
diff --git a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/fsck/ReconReplicationManager.java b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/fsck/ReconReplicationManager.java
index 09005cadb18e..af7b3045e068 100644
--- a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/fsck/ReconReplicationManager.java
+++ b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/fsck/ReconReplicationManager.java
@@ -17,6 +17,10 @@
package org.apache.hadoop.ozone.recon.fsck;
+import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState.CLOSED;
+import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType.RATIS;
+import static org.apache.hadoop.hdds.scm.container.ContainerReplicaChecksumMismatch.hasMismatch;
+
import java.io.IOException;
import java.time.Clock;
import java.util.ArrayList;
@@ -24,7 +28,6 @@
import java.util.HashSet;
import java.util.List;
import java.util.Map;
-import java.util.Objects;
import java.util.Set;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.scm.PlacementPolicy;
@@ -226,43 +229,6 @@ public synchronized void start() {
// Do nothing - we call processAll() manually from ContainerHealthTask
}
- /**
- * Checks if container replicas have mismatched data checksums.
- * This is a Recon-specific check not done by SCM's ReplicationManager.
- *
- * REPLICA_MISMATCH detection is crucial for identifying:
- *
- * - Bit rot (silent data corruption)
- * - Failed writes to some replicas
- * - Storage corruption on specific datanodes
- * - Network corruption during replication
- *
- *
- *
- * This uses checksum mismatch logic:
- * {@code replicas.stream().map(ContainerReplica::getDataChecksum).distinct().count() != 1}
- *
- *
- * @param replicas Set of container replicas to check
- * @return true if replicas have different data checksums
- */
- private boolean hasDataChecksumMismatch(Set replicas) {
- if (replicas == null || replicas.isEmpty()) {
- return false;
- }
-
- // Count distinct checksums (filter out nulls)
- long distinctChecksums = replicas.stream()
- .map(ContainerReplica::getDataChecksum)
- .filter(Objects::nonNull)
- .distinct()
- .count();
-
- // More than 1 distinct checksum = data mismatch
- // 0 distinct checksums = all nulls, no mismatch
- return distinctChecksums > 1;
- }
-
/**
* Override processAll() to capture ALL per-container health states,
* not just aggregate counts and 100 samples.
@@ -271,7 +237,7 @@ private boolean hasDataChecksumMismatch(Set replicas) {
*
* - Get all containers from ContainerManager
* - Process each container using inherited health check chain (SCM logic)
- * - Additionally check for REPLICA_MISMATCH (Recon-specific)
+ * - Persist checksum mismatches as Recon REPLICA_MISMATCH records
* - Capture ALL unhealthy container IDs per health state (no sampling limit)
* - Store results in Recon's UNHEALTHY_CONTAINERS table
*
@@ -280,7 +246,7 @@ private boolean hasDataChecksumMismatch(Set replicas) {
*
* - Uses ReconReplicationManagerReport (captures all containers)
* - Uses MonitoringReplicationQueue (doesn't enqueue commands)
- * - Adds REPLICA_MISMATCH detection (not done by SCM)
+ * - Persists REPLICA_MISMATCH health records
* - Stores results in database instead of just keeping in-memory report
*
*/
@@ -313,8 +279,11 @@ public synchronized void processAll() {
// readOnly=true ensures no commands are generated
processContainer(container, replicas, pendingOps, nullQueue, report, true);
- // ADDITIONAL CHECK: Detect REPLICA_MISMATCH (Recon-specific, not in SCM)
- if (hasDataChecksumMismatch(replicas)) {
+ // Persist checksum mismatches in Recon's REPLICA_MISMATCH state.
+ if (container.getState() == CLOSED &&
+ container.getReplicationType() == RATIS &&
+ hasMismatch(replicas, ContainerReplica::getSequenceId,
+ ContainerReplica::getDataChecksum)) {
report.addReplicaMismatchContainer(cid);
LOG.debug("Container {} has data checksum mismatch across replicas", cid);
}
diff --git a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/fsck/TestReconReplicationManager.java b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/fsck/TestReconReplicationManager.java
index 0678caa86eb7..d1115299f8c7 100644
--- a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/fsck/TestReconReplicationManager.java
+++ b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/fsck/TestReconReplicationManager.java
@@ -419,6 +419,8 @@ private ContainerInfo mockContainerInfo(long containerId, long numberOfKeys,
when(containerInfo.getNumberOfKeys()).thenReturn(numberOfKeys);
when(containerInfo.getUsedBytes()).thenReturn(usedBytes);
when(containerInfo.getReplicationConfig()).thenReturn(replicationConfig);
+ when(containerInfo.getReplicationType()).thenReturn(
+ HddsProtos.ReplicationType.RATIS);
when(containerInfo.getState()).thenReturn(HddsProtos.LifeCycleState.CLOSED);
when(containerInfo.getHealthState()).thenAnswer(invocation -> healthStateRef.get());
doAnswer(invocation -> {
@@ -427,6 +429,8 @@ private ContainerInfo mockContainerInfo(long containerId, long numberOfKeys,
}).when(containerInfo).setHealthState(
org.mockito.ArgumentMatchers.any(ContainerHealthState.class));
when(replicationConfig.getRequiredNodes()).thenReturn(requiredNodes);
+ when(replicationConfig.getReplicationType()).thenReturn(
+ HddsProtos.ReplicationType.RATIS);
return containerInfo;
}
@@ -502,6 +506,7 @@ private Set setOfMockReplicasWithChecksums(Long... checksums)
for (Long checksum : checksums) {
ContainerReplica replica = mock(ContainerReplica.class);
when(replica.getDataChecksum()).thenReturn(checksum);
+ when(replica.getSequenceId()).thenReturn(1L);
replicas.add(replica);
}
return replicas;
From fa12e8845b8c7c92545352af7021744e281a47f3 Mon Sep 17 00:00:00 2001
From: f64116045
Date: Sat, 5 Sep 2026 22:54:37 +0800
Subject: [PATCH 2/2] HDDS-16012. Fix Recon checksum mismatch test
---
.../hadoop/ozone/recon/TestReconTasks.java | 36 +++++++++++--------
1 file changed, 21 insertions(+), 15 deletions(-)
diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java
index 4d9f802aee3a..fce9d0cddcb8 100644
--- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java
+++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconTasks.java
@@ -670,13 +670,11 @@ public void testContainerHealthTaskDetectsOverReplicatedAndNegativeSize()
* when replicas of a CLOSED RF3 container report different data checksums, and
* that the state clears once the checksums are made uniform again.
*
- * Strategy: after writing data and closing the RF3 container, one replica's
- * checksum is replaced with a non-zero value via {@link ContainerReplica#toBuilder()}.
- * Because {@link ContainerReplica} is immutable,
- * {@code reconCm.updateContainerReplica()} replaces the existing replica entry
- * (keyed by {@code containerID + datanodeDetails}) with the modified copy.
- * This makes the set have distinct checksums (e.g. {0, 0, 12345}), which triggers
- * {@code hasDataChecksumMismatch()}'s {@code distinctChecksums > 1} check.
+ * Strategy: after writing data and closing the RF3 container, all replicas
+ * are given the same non-zero sequence ID and checksum, then one replica's
+ * checksum is changed via {@link ContainerReplica#toBuilder()}. Because
+ * {@link ContainerReplica} is immutable, {@code reconCm.updateContainerReplica()}
+ * replaces the existing replica entry keyed by container ID and datanode.
*
* Note: {@code REPLICA_MISMATCH} records are now properly cleaned up by
* {@code batchDeleteSCMStatesForContainers} on each scan cycle (previously they
@@ -729,15 +727,23 @@ public void testContainerHealthTaskDetectsReplicaMismatch() throws Exception {
reconCm.updateContainerState(
containerInfo.containerID(), HddsProtos.LifeCycleEvent.CLOSE);
- // Inject a checksum mismatch: replace one replica with an identical copy
- // that has a non-zero dataChecksum. The other two replicas have checksum=0
- // (ContainerChecksums.unknown()), giving distinct checksums {0, 12345} and
- // triggering distinctChecksums > 1.
+ // Give every replica a known checksum first. A zero checksum means it has
+ // not been reported yet and must not be treated as a mismatch.
Set currentReplicas = reconCm.getContainerReplicas(cid);
- ContainerReplica originalReplica = currentReplicas.iterator().next();
- ContainerReplica mismatchedReplica = originalReplica.toBuilder()
+ for (ContainerReplica replica : currentReplicas) {
+ reconCm.updateContainerReplica(cid, replica.toBuilder()
+ .setSequenceId(1L)
+ .setChecksums(ContainerChecksums.of(12345L))
+ .build());
+ }
+
+ ContainerReplica matchingReplica = currentReplicas.iterator().next().toBuilder()
+ .setSequenceId(1L)
.setChecksums(ContainerChecksums.of(12345L))
.build();
+ ContainerReplica mismatchedReplica = matchingReplica.toBuilder()
+ .setChecksums(ContainerChecksums.of(54321L))
+ .build();
reconCm.updateContainerReplica(cid, mismatchedReplica);
forceContainerHealthScan(reconScm);
@@ -749,8 +755,8 @@ public void testContainerHealthTaskDetectsReplicaMismatch() throws Exception {
assertTrue(containsContainerId(replicaMismatch, containerID),
"Container with differing replica checksums should be REPLICA_MISMATCH");
- // Recovery: restore the original replica (uniform checksums → no mismatch).
- reconCm.updateContainerReplica(cid, originalReplica);
+ // Recovery: restore the matching checksum.
+ reconCm.updateContainerReplica(cid, matchingReplica);
forceContainerHealthScan(reconScm);
List replicaMismatchAfterRecovery =