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 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/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 = 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) { *
    *
  1. Get all containers from ContainerManager
  2. *
  3. Process each container using inherited health check chain (SCM logic)
  4. - *
  5. Additionally check for REPLICA_MISMATCH (Recon-specific)
  6. + *
  7. Persist checksum mismatches as Recon REPLICA_MISMATCH records
  8. *
  9. Capture ALL unhealthy container IDs per health state (no sampling limit)
  10. *
  11. Store results in Recon's UNHEALTHY_CONTAINERS table
  12. *
@@ -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;