diff --git a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto index 6c39fc22cc4f..8aa022e74d43 100644 --- a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto +++ b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto @@ -675,6 +675,19 @@ message ContainerBalancerTaskIterationStatusInfoProto { repeated NodeTransferInfoProto sizeEnteringNodes = 9; repeated NodeTransferInfoProto sizeLeavingNodes = 10; optional int64 iterationDuration = 11; + repeated ContainerMoveFailureDetailProto containerMoveFailures = 12; +} + +message ContainerMoveFailureDetailProto { + optional string reason = 1; + optional int64 count = 2; + repeated NodeFailureCountProto sourceFailureCounts = 3; + repeated NodeFailureCountProto targetFailureCounts = 4; +} + +message NodeFailureCountProto { + optional string datanodeUuid = 1; + optional int64 count = 2; } message NodeTransferInfoProto { diff --git a/hadoop-hdds/interface-admin/src/main/resources/proto.lock b/hadoop-hdds/interface-admin/src/main/resources/proto.lock index 02184011b692..0ba44234369d 100644 --- a/hadoop-hdds/interface-admin/src/main/resources/proto.lock +++ b/hadoop-hdds/interface-admin/src/main/resources/proto.lock @@ -2339,6 +2339,58 @@ "name": "iterationDuration", "type": "int64", "optional": true + }, + { + "id": 12, + "name": "containerMoveFailures", + "type": "ContainerMoveFailureDetailProto", + "is_repeated": true + } + ] + }, + { + "name": "ContainerMoveFailureDetailProto", + "fields": [ + { + "id": 1, + "name": "reason", + "type": "string", + "optional": true + }, + { + "id": 2, + "name": "count", + "type": "int64", + "optional": true + }, + { + "id": 3, + "name": "sourceFailureCounts", + "type": "NodeFailureCountProto", + "is_repeated": true + }, + { + "id": 4, + "name": "targetFailureCounts", + "type": "NodeFailureCountProto", + "is_repeated": true + } + ] + }, + { + "name": "NodeFailureCountProto", + "fields": [ + { + "id": 1, + "name": "datanodeUuid", + "type": "string", + "optional": true + }, + { + "id": 2, + "name": "count", + "type": "int64", + "optional": true } ] }, diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java index bb2b4c873254..1c14ed1b124f 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTask.java @@ -37,6 +37,7 @@ import java.util.Objects; import java.util.Queue; import java.util.Set; +import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; @@ -94,6 +95,7 @@ public class ContainerBalancerTask implements Runnable { private Set includeNodes; private ContainerBalancerConfiguration config; private ContainerBalancerMetrics metrics; + private final ContainerMoveFailureTracker moveFailureTracker; private double upperLimit; private double lowerLimit; private ContainerBalancerSelectionCriteria selectionCriteria; @@ -151,6 +153,7 @@ public ContainerBalancerTask(StorageContainerManager scm, this.containerBalancer = containerBalancer; this.config = config; this.metrics = metrics; + this.moveFailureTracker = new ContainerMoveFailureTracker(); this.scmContext = scm.getScmContext(); this.overUtilizedNodes = new ArrayList<>(); this.underUtilizedNodes = new ArrayList<>(); @@ -346,7 +349,7 @@ private ContainerBalancerTaskIterationStatusInfo getIterationStatistic(Integer i currentIterationResultName, iterationDuration ); - ContainerMoveInfo containerMoveInfo = new ContainerMoveInfo(metrics); + ContainerMoveInfo containerMoveInfo = new ContainerMoveInfo(metrics, moveFailureTracker); DataMoveInfo dataMoveInfo = getDataMoveInfo(sizeEnteringDataToNodes, sizeLeavingDataFromNodes); @@ -793,7 +796,8 @@ private void checkIterationMoveResults() { LOG.warn("Container balancer is interrupted"); Thread.currentThread().interrupt(); } catch (TimeoutException e) { - long timeoutCounts = cancelMovesThatExceedTimeoutDuration(); + long timeoutCounts = cancelMovesThatExceedTimeoutDuration( + ContainerMoveFailureReason.ITERATION_MOVE_TIMEOUT.name()); LOG.warn("{} Container moves are canceled.", timeoutCounts); metrics.incrementNumContainerMovesTimeoutInLatestIteration( timeoutCounts); @@ -836,7 +840,7 @@ private void checkIterationMoveResults() { * @return number of moves that did not complete (timed out) and were * cancelled. */ - private long cancelMovesThatExceedTimeoutDuration() { + private long cancelMovesThatExceedTimeoutDuration(String reason) { Set>> entries = moveSelectionToFutureMap.entrySet(); @@ -851,11 +855,13 @@ private long cancelMovesThatExceedTimeoutDuration() { CompletableFuture> entry = iterator.next(); if (!entry.getValue().isDone()) { + ContainerMoveSelection moveSelection = entry.getKey(); + ContainerID containerID = moveSelection.getContainerID(); + DatanodeDetails source = containerToSourceMap.get(containerID); + DatanodeDetails target = moveSelection.getTargetNode(); LOG.warn("Container move timed out for container {} from source {}" + - " to target {}.", entry.getKey().getContainerID(), - containerToSourceMap.get(entry.getKey().getContainerID()), - entry.getKey().getTargetNode()); - + " to target {}.", containerID, source, target); + moveFailureTracker.recordFailure(reason, source, target); entry.getValue().cancel(true); numCancelled += 1; } @@ -1004,11 +1010,15 @@ private boolean moveContainer(DatanodeDetails source, metrics.incrementCurrentIterationContainerMoveMetric(result, 1); moveSelectionToFutureMap.remove(moveSelection); if (ex != null) { - LOG.info("Container move for container {} from source {} to " + - "target {} failed with exceptions.", - containerID, source, - moveSelection.getTargetNode(), ex); - metrics.incrementNumContainerMovesFailedInLatestIteration(1); + if (!isCancellationCause(ex)) { + LOG.info("Container move for container {} from source {} to " + + "target {} failed with exceptions.", + containerID, source, + moveSelection.getTargetNode(), ex); + metrics.incrementNumContainerMovesFailedInLatestIteration(1); + moveFailureTracker.recordFailure(MoveManager.MoveResult.FAIL_UNEXPECTED_ERROR, + source, moveSelection.getTargetNode()); + } } else { if (result == MoveManager.MoveResult.COMPLETED) { metrics.incrementDataSizeMovedInLatestIteration(containerInfo.getUsedBytes()); @@ -1021,6 +1031,7 @@ private boolean moveContainer(DatanodeDetails source, " {} failed: {}", moveSelection.getContainerID(), source, moveSelection.getTargetNode(), result); + moveFailureTracker.recordFailure(result, source, moveSelection.getTargetNode()); } } }); @@ -1032,14 +1043,20 @@ private boolean moveContainer(DatanodeDetails source, // exclude the container which caused failure of move to avoid error in next run. selectionCriteria.addToExcludeDueToFailContainers(moveSelection.getContainerID()); metrics.incrementNumContainerMovesFailedInLatestIteration(1); + moveFailureTracker.recordFailure(ContainerMoveFailureReason.PRE_MOVE_CONTAINER_NOT_FOUND.name(), + source, moveSelection.getTargetNode()); return false; } catch (NodeNotFoundException e) { LOG.warn("Container move failed for container {}", containerID, e); metrics.incrementNumContainerMovesFailedInLatestIteration(1); + moveFailureTracker.recordFailure(ContainerMoveFailureReason.PRE_MOVE_NODE_NOT_FOUND.name(), + source, moveSelection.getTargetNode()); return false; } catch (ContainerReplicaNotFoundException e) { LOG.warn("Container move failed for container {}", containerID, e); metrics.incrementNumContainerMovesFailedInLatestIteration(1); + moveFailureTracker.recordFailure(ContainerMoveFailureReason.PRE_MOVE_REPLICA_NOT_FOUND.name(), + source, moveSelection.getTargetNode()); // add source back to queue for replica not found only // the container is not excluded as it is a replica related failure findSourceStrategy.addBackSourceDataNode(source); @@ -1225,6 +1242,11 @@ private void resetState() { metrics.resetDataSizeUnbalancedGB(); metrics.resetNumDatanodesUnbalanced(); metrics.resetNumContainerMovesFailedInLatestIteration(); + moveFailureTracker.reset(); + } + + private static boolean isCancellationCause(Throwable ex) { + return ex instanceof CancellationException; } /** @@ -1342,4 +1364,18 @@ enum Status { STOPPING, STOPPED } + + /** + * Recorded when a scheduled container move fails before {@link MoveManager#move} + * completes, or when the balancer stops waiting for in-flight moves. + */ + enum ContainerMoveFailureReason { + PRE_MOVE_CONTAINER_NOT_FOUND, + PRE_MOVE_NODE_NOT_FOUND, + PRE_MOVE_REPLICA_NOT_FOUND, + /** + * Move did not complete before the iteration move wait timeout expired. + */ + ITERATION_MOVE_TIMEOUT + } } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTaskIterationStatusInfo.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTaskIterationStatusInfo.java index f16f3cfe7e54..79875f441031 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTaskIterationStatusInfo.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerBalancerTaskIterationStatusInfo.java @@ -106,6 +106,14 @@ public long getContainerMovesTimeout() { return containerMoveInfo.getContainerMovesTimeout(); } + /** + * Get per-reason failure summaries with per-datanode counts for this iteration. + * @return list of failure details, one entry per failure reason + */ + public List getFailures() { + return containerMoveInfo.getFailures(); + } + /** * Get a map of the node IDs and the corresponding data sizes moved to each node. * @return nodeId to size entering from node map @@ -151,9 +159,35 @@ public StorageContainerLocationProtocolProtos.ContainerBalancerTaskIterationStat .addAllSizeLeavingNodes( mapToProtoNodeTransferInfo(getSizeLeavingNodes()) ) + .addAllContainerMoveFailures(mapToProtoFailures(getFailures())) .build(); } + private List mapToProtoFailures( + List failures) { + return failures.stream() + .map(failure -> { + StorageContainerLocationProtocolProtos.ContainerMoveFailureDetailProto.Builder builder = + StorageContainerLocationProtocolProtos.ContainerMoveFailureDetailProto.newBuilder() + .setReason(failure.getReason()) + .setCount(failure.getCount()); + failure.getSourceFailureCounts().forEach((uuid, count) -> + builder.addSourceFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid(uuid) + .setCount(count) + .build())); + failure.getTargetFailureCounts().forEach((uuid, count) -> + builder.addTargetFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid(uuid) + .setCount(count) + .build())); + return builder.build(); + }) + .collect(Collectors.toList()); + } + /** * Converts an instance into the protobuf compatible object. * @param nodes node id to node traffic size diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveFailureDetail.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveFailureDetail.java new file mode 100644 index 000000000000..d890a80bd0cf --- /dev/null +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveFailureDetail.java @@ -0,0 +1,55 @@ +/* + * 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.balancer; + +import java.util.Collections; +import java.util.Map; + +/** + * Per-reason container move failure summary with per-datanode failure counts. + */ +public final class ContainerMoveFailureDetail { + private final String reason; + private final long count; + private final Map sourceFailureCounts; + private final Map targetFailureCounts; + + public ContainerMoveFailureDetail(String reason, long count, + Map sourceFailureCounts, Map targetFailureCounts) { + this.reason = reason; + this.count = count; + this.sourceFailureCounts = Collections.unmodifiableMap(sourceFailureCounts); + this.targetFailureCounts = Collections.unmodifiableMap(targetFailureCounts); + } + + public String getReason() { + return reason; + } + + public long getCount() { + return count; + } + + public Map getSourceFailureCounts() { + return sourceFailureCounts; + } + + public Map getTargetFailureCounts() { + return targetFailureCounts; + } +} diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveFailureTracker.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveFailureTracker.java new file mode 100644 index 000000000000..9d5ed63c9443 --- /dev/null +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveFailureTracker.java @@ -0,0 +1,69 @@ +/* + * 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.balancer; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; + +/** + * Tracks per-iteration container move failures by reason and per-datanode counts. + */ +public final class ContainerMoveFailureTracker { + private final Map failuresByReason = new HashMap<>(); + private final Map> sourceFailureCountsByReason = new HashMap<>(); + private final Map> targetFailureCountsByReason = new HashMap<>(); + + public synchronized void recordFailure(MoveManager.MoveResult result, DatanodeDetails source, + DatanodeDetails target) { + recordFailure(result.name(), source, target); + } + + public synchronized void recordFailure(String reason, DatanodeDetails source, DatanodeDetails target) { + failuresByReason.merge(reason, 1L, Long::sum); + if (source != null) { + sourceFailureCountsByReason.computeIfAbsent(reason, k -> new HashMap<>()) + .merge(source.getUuidString(), 1L, Long::sum); + } + if (target != null) { + targetFailureCountsByReason.computeIfAbsent(reason, k -> new HashMap<>()) + .merge(target.getUuidString(), 1L, Long::sum); + } + } + + public synchronized void reset() { + failuresByReason.clear(); + sourceFailureCountsByReason.clear(); + targetFailureCountsByReason.clear(); + } + + public synchronized List getFailures() { + List result = new ArrayList<>(); + for (Map.Entry entry : failuresByReason.entrySet()) { + String reason = entry.getKey(); + long count = entry.getValue(); + Map srcCounts = sourceFailureCountsByReason.getOrDefault(reason, Collections.emptyMap()); + Map tgtCounts = targetFailureCountsByReason.getOrDefault(reason, Collections.emptyMap()); + result.add(new ContainerMoveFailureDetail(reason, count, new HashMap<>(srcCounts), new HashMap<>(tgtCounts))); + } + return result; + } +} diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveInfo.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveInfo.java index 3b3c7f91933e..2dab77120ddb 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveInfo.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/balancer/ContainerMoveInfo.java @@ -17,6 +17,8 @@ package org.apache.hadoop.hdds.scm.container.balancer; +import java.util.List; + /** * Information about moving containers. */ @@ -25,20 +27,23 @@ public class ContainerMoveInfo { private final long containerMovesCompleted; private final long containerMovesFailed; private final long containerMovesTimeout; + private final List failures; public ContainerMoveInfo(long containerMovesScheduled, long containerMovesCompleted, long containerMovesFailed, - long containerMovesTimeout) { + long containerMovesTimeout, List failures) { this.containerMovesScheduled = containerMovesScheduled; this.containerMovesCompleted = containerMovesCompleted; this.containerMovesFailed = containerMovesFailed; this.containerMovesTimeout = containerMovesTimeout; + this.failures = failures; } - public ContainerMoveInfo(ContainerBalancerMetrics metrics) { + public ContainerMoveInfo(ContainerBalancerMetrics metrics, ContainerMoveFailureTracker failureTracker) { this.containerMovesScheduled = metrics.getNumContainerMovesScheduledInLatestIteration(); this.containerMovesCompleted = metrics.getNumContainerMovesCompletedInLatestIteration(); this.containerMovesFailed = metrics.getNumContainerMovesFailedInLatestIteration(); this.containerMovesTimeout = metrics.getNumContainerMovesTimeoutInLatestIteration(); + this.failures = failureTracker.getFailures(); } public long getContainerMovesScheduled() { @@ -56,4 +61,8 @@ public long getContainerMovesFailed() { public long getContainerMovesTimeout() { return containerMovesTimeout; } + + public List getFailures() { + return failures; + } } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerMoveFailureBreakdown.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerMoveFailureBreakdown.java new file mode 100644 index 000000000000..3b028168270e --- /dev/null +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/balancer/TestContainerBalancerMoveFailureBreakdown.java @@ -0,0 +1,248 @@ +/* + * 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.balancer; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.atLeast; +import static org.mockito.Mockito.atLeastOnce; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.time.Duration; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerNotFoundException; +import org.apache.hadoop.hdds.scm.container.ContainerReplicaNotFoundException; +import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException; +import org.apache.hadoop.ozone.OzoneConsts; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +/** + * Tests that {@link ContainerBalancerTask} records failure breakdown and details + * in iteration statistics for common move failure reasons. + */ +class TestContainerBalancerMoveFailureBreakdown { + + private static final int NODE_COUNT = 5; + private static final long STORAGE_UNIT = OzoneConsts.GB; + + @Test + void testReplicationTimeoutRecordedInFailureBreakdown() + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + MockedSCM mockedScm = createMockedScm(); + mockMoveFailureOnce(mockedScm, MoveManager.MoveResult.REPLICATION_FAIL_TIME_OUT); + + ContainerBalancerTask task = mockedScm.startBalancerTask(buildConfig(mockedScm)); + ContainerBalancerTaskIterationStatusInfo iteration = getCompletedIteration(task); + + assertThat(iteration.getContainerMovesTimeout()).isEqualTo(1); + assertFailureBreakdown(mockedScm, iteration, MoveManager.MoveResult.REPLICATION_FAIL_TIME_OUT.name(), 0); + } + + @Test + void testReplicationNotHealthyAfterMoveRecordedInFailureBreakdown() + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + MockedSCM mockedScm = createMockedScm(); + mockMoveFailureOnce(mockedScm, MoveManager.MoveResult.REPLICATION_NOT_HEALTHY_AFTER_MOVE); + + ContainerBalancerTask task = mockedScm.startBalancerTask(buildConfig(mockedScm)); + ContainerBalancerTaskIterationStatusInfo iteration = getCompletedIteration(task); + + assertThat(iteration.getContainerMovesFailed()).isEqualTo(1); + assertFailureBreakdown(mockedScm, iteration, + MoveManager.MoveResult.REPLICATION_NOT_HEALTHY_AFTER_MOVE.name(), 0); + } + + @Test + void testDeletionTimeoutRecordedInFailureBreakdown() + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + MockedSCM mockedScm = createMockedScm(); + mockMoveFailureOnce(mockedScm, MoveManager.MoveResult.DELETION_FAIL_TIME_OUT); + + ContainerBalancerTask task = mockedScm.startBalancerTask(buildConfig(mockedScm)); + ContainerBalancerTaskIterationStatusInfo iteration = getCompletedIteration(task); + + assertThat(iteration.getContainerMovesTimeout()).isEqualTo(1); + assertFailureBreakdown(mockedScm, iteration, MoveManager.MoveResult.DELETION_FAIL_TIME_OUT.name(), 0); + } + + @Test + void testIterationMoveTimeoutRecordedInFailureBreakdown() + throws NodeNotFoundException, ContainerNotFoundException, ContainerReplicaNotFoundException { + MockedSCM mockedScm = createMockedScm(); + ContainerBalancerConfiguration config = buildConfig(mockedScm); + config.setMoveTimeout(Duration.ofMillis(50)); + mockFirstMoveNeverCompletes(mockedScm); + + ContainerBalancerTask task = mockedScm.startBalancerTask(config); + ContainerBalancerTaskIterationStatusInfo iteration = getCompletedIteration(task); + + assertThat(iteration.getContainerMovesTimeout()).isEqualTo(1); + assertFailureBreakdown(mockedScm, iteration, + ContainerBalancerTask.ContainerMoveFailureReason.ITERATION_MOVE_TIMEOUT.name(), 0); + } + + @Test + void testFailureBreakdownTotalsMatchHeadlineCountersInControlledScenario() + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + MockedSCM mockedScm = createMockedScm(); + when(mockedScm.getMoveManager().move(any(ContainerID.class), + any(DatanodeDetails.class), any(DatanodeDetails.class))) + .thenReturn(CompletableFuture.completedFuture( + MoveManager.MoveResult.REPLICATION_FAIL_TIME_OUT)) + .thenReturn(CompletableFuture.completedFuture( + MoveManager.MoveResult.REPLICATION_NOT_HEALTHY_AFTER_MOVE)) + .thenReturn(CompletableFuture.completedFuture(MoveManager.MoveResult.COMPLETED)); + + ContainerBalancerTask task = mockedScm.startBalancerTask(buildConfig(mockedScm)); + ContainerBalancerTaskIterationStatusInfo iteration = getCompletedIteration(task); + + assertThat(iteration.getContainerMovesTimeout()).isEqualTo(1); + assertThat(iteration.getContainerMovesFailed()).isEqualTo(1); + verify(mockedScm.getMoveManager(), atLeast(2)).move( + any(ContainerID.class), any(DatanodeDetails.class), any(DatanodeDetails.class)); + assertBreakdownTotalsMatchHeadlineCounters(iteration); + } + + @Test + void testPreMoveContainerNotFoundRecordedInFailureBreakdown() + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + MockedSCM mockedScm = createMockedScm(); + when(mockedScm.getMoveManager().move(any(ContainerID.class), + any(DatanodeDetails.class), any(DatanodeDetails.class))) + .thenThrow(ContainerNotFoundException.newInstanceForTesting()) + .thenReturn(CompletableFuture.completedFuture(MoveManager.MoveResult.COMPLETED)); + + ContainerBalancerTask task = mockedScm.startBalancerTask(buildConfig(mockedScm)); + ContainerBalancerTaskIterationStatusInfo iteration = getCompletedIteration(task); + + assertThat(iteration.getContainerMovesFailed()).isEqualTo(1); + assertFailureBreakdown(mockedScm, iteration, + ContainerBalancerTask.ContainerMoveFailureReason.PRE_MOVE_CONTAINER_NOT_FOUND.name(), 0); + } + + @Test + void testSameReasonAggregatesSourceFailureCountsFromMultipleSources() { + ContainerMoveFailureTracker tracker = new ContainerMoveFailureTracker(); + DatanodeDetails source1 = MockDatanodeDetails.createDatanodeDetails("1.1.1.1", "/r1"); + DatanodeDetails source2 = MockDatanodeDetails.createDatanodeDetails("2.2.2.2", "/r1"); + DatanodeDetails target = MockDatanodeDetails.createDatanodeDetails("3.3.3.3", "/r2"); + + String reason = MoveManager.MoveResult.REPLICATION_FAIL_TIME_OUT.name(); + tracker.recordFailure(reason, source1, target); + tracker.recordFailure(reason, source2, target); + + ContainerMoveFailureDetail detail = tracker.getFailures().stream() + .filter(f -> reason.equals(f.getReason())) + .findFirst() + .orElse(null); + assertThat(detail).as("failure detail for reason " + reason).isNotNull(); + assertThat(detail.getCount()).isEqualTo(2L); + assertThat(detail.getSourceFailureCounts()) + .hasSize(2) + .containsEntry(source1.getUuidString(), 1L) + .containsEntry(source2.getUuidString(), 1L); + assertThat(detail.getTargetFailureCounts()) + .hasSize(1) + .containsEntry(target.getUuidString(), 2L); + } + + private static MockedSCM createMockedScm() { + return new MockedSCM(new MockCluster(NODE_COUNT, STORAGE_UNIT)); + } + + private static void mockMoveFailureOnce(MockedSCM mockedScm, MoveManager.MoveResult failureResult) + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + when(mockedScm.getMoveManager().move(any(ContainerID.class), + any(DatanodeDetails.class), any(DatanodeDetails.class))) + .thenReturn(CompletableFuture.completedFuture(failureResult)) + .thenReturn(CompletableFuture.completedFuture(MoveManager.MoveResult.COMPLETED)); + } + + private static void mockFirstMoveNeverCompletes(MockedSCM mockedScm) + throws NodeNotFoundException, ContainerReplicaNotFoundException, ContainerNotFoundException { + AtomicInteger moveInvocations = new AtomicInteger(0); + when(mockedScm.getMoveManager().move(any(ContainerID.class), + any(DatanodeDetails.class), any(DatanodeDetails.class))) + .thenAnswer(invocation -> { + if (moveInvocations.getAndIncrement() == 0) { + return new CompletableFuture<>(); + } + return CompletableFuture.completedFuture(MoveManager.MoveResult.COMPLETED); + }); + } + + private static ContainerBalancerConfiguration buildConfig(MockedSCM mockedScm) { + ContainerBalancerConfiguration config = new ContainerBalancerConfigBuilder(mockedScm.getNodeCount()).build(); + config.setMaxSizeToMovePerIteration(5 * STORAGE_UNIT); + config.setMaxSizeEnteringTarget(5 * STORAGE_UNIT); + config.setMaxDatanodesPercentageToInvolvePerIteration(100); + return config; + } + + private static ContainerBalancerTaskIterationStatusInfo getCompletedIteration( + ContainerBalancerTask task) { + List iterations = + task.getCurrentIterationsStatistic(); + assertEquals(1, iterations.size()); + ContainerBalancerTaskIterationStatusInfo iteration = iterations.get(0); + assertEquals("ITERATION_COMPLETED", iteration.getIterationResult()); + return iteration; + } + + private static void assertBreakdownTotalsMatchHeadlineCounters( + ContainerBalancerTaskIterationStatusInfo iteration) { + long breakdownTotal = iteration.getFailures().stream() + .mapToLong(ContainerMoveFailureDetail::getCount) + .sum(); + long headlineTotal = iteration.getContainerMovesFailed() + iteration.getContainerMovesTimeout(); + assertEquals(headlineTotal, breakdownTotal, + "sum(failure breakdown) should equal failed + timeout counters"); + } + + private static void assertFailureBreakdown(MockedSCM mockedScm, ContainerBalancerTaskIterationStatusInfo iteration, + String expectedReason, int moveIndex) throws NodeNotFoundException, ContainerReplicaNotFoundException, + ContainerNotFoundException { + ArgumentCaptor sourceCaptor = ArgumentCaptor.forClass(DatanodeDetails.class); + ArgumentCaptor targetCaptor = ArgumentCaptor.forClass(DatanodeDetails.class); + verify(mockedScm.getMoveManager(), atLeastOnce()).move( + any(ContainerID.class), sourceCaptor.capture(), targetCaptor.capture()); + assertThat(sourceCaptor.getAllValues().size()).isGreaterThan(moveIndex); + assertThat(targetCaptor.getAllValues().size()).isGreaterThan(moveIndex); + String sourceUuid = sourceCaptor.getAllValues().get(moveIndex).getUuidString(); + String targetUuid = targetCaptor.getAllValues().get(moveIndex).getUuidString(); + + List failures = iteration.getFailures(); + assertThat(failures).as("failure details").isNotEmpty(); + ContainerMoveFailureDetail detail = failures.stream() + .filter(f -> expectedReason.equals(f.getReason())) + .findFirst() + .orElse(null); + assertThat(detail).as("failure detail for reason " + expectedReason).isNotNull(); + assertThat(detail.getCount()).isEqualTo(1L); + assertThat(detail.getSourceFailureCounts()).hasSize(1).containsEntry(sourceUuid, 1L); + assertThat(detail.getTargetFailureCounts()).hasSize(1).containsEntry(targetUuid, 1L); + } +} diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerStatusSubcommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerStatusSubcommand.java index a55fc7906dba..c23bd1a063a1 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerStatusSubcommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerBalancerStatusSubcommand.java @@ -27,6 +27,7 @@ import java.time.OffsetDateTime; import java.time.ZoneId; import java.time.format.DateTimeFormatter; +import java.util.Comparator; import java.util.List; import java.util.stream.Collectors; import org.apache.hadoop.hdds.cli.HddsVersionProvider; @@ -34,6 +35,8 @@ import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerStatusInfoResponseProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerBalancerTaskIterationStatusInfoProto; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerMoveFailureDetailProto; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.NodeFailureCountProto; import org.apache.hadoop.hdds.scm.client.ScmClient; import org.apache.hadoop.ozone.OzoneConsts; import picocli.CommandLine; @@ -232,6 +235,7 @@ private String getPrettyIterationStatusInfo(ContainerBalancerTaskIterationStatus if (leavingDataNodeList.isEmpty()) { leavingDataNodeList = " -" + System.lineSeparator(); } + String failures = formatFailures(iterationStatusInfo.getContainerMoveFailuresList()); return String.format( "%-50s %s%n" + "%-50s %s%n" + @@ -243,6 +247,7 @@ private String getPrettyIterationStatusInfo(ContainerBalancerTaskIterationStatus "%-50s %s%n" + "%-50s %s%n" + "%-50s %s%n" + + "%s" + "%-50s %n%s" + "%-50s %n%s", "Key", "Value", @@ -256,9 +261,38 @@ private String getPrettyIterationStatusInfo(ContainerBalancerTaskIterationStatus "Already moved containers", containerMovesCompleted, "Failed to move containers", containerMovesFailed, "Failed to move containers by timeout", containerMovesTimeout, + failures, "Entered data to nodes", enteringDataNodeList, "Exited data from nodes", leavingDataNodeList); } + private String formatFailures(List failures) { + if (failures.isEmpty()) { + return ""; + } + List sorted = failures.stream() + .sorted(Comparator.comparingLong(ContainerMoveFailureDetailProto::getCount).reversed() + .thenComparing(ContainerMoveFailureDetailProto::getReason)) + .collect(Collectors.toList()); + StringBuilder builder = new StringBuilder(); + builder.append(String.format("%-50s %n", "Failed container moves")); + for (ContainerMoveFailureDetailProto failure : sorted) { + builder.append(String.format(" %-48s %d%n", failure.getReason(), failure.getCount())); + if (!failure.getSourceFailureCountsList().isEmpty()) { + builder.append(String.format(" %-46s %n", "Source datanodes")); + for (NodeFailureCountProto src : failure.getSourceFailureCountsList()) { + builder.append(String.format(" %-44s %d%n", src.getDatanodeUuid(), src.getCount())); + } + } + if (!failure.getTargetFailureCountsList().isEmpty()) { + builder.append(String.format(" %-46s %n", "Target datanodes")); + for (NodeFailureCountProto tgt : failure.getTargetFailureCountsList()) { + builder.append(String.format(" %-44s %d%n", tgt.getDatanodeUuid(), tgt.getCount())); + } + } + } + return builder.toString(); + } + } diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java index 407785d48a33..f00f76da6f78 100644 --- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java +++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestContainerBalancerSubCommand.java @@ -712,4 +712,78 @@ void testContainerBalancerStatusSubcommandStoppedAfterAllIterationsCompleteVerbo .contains(ITERATION_2_COMPLETED_OUTPUT) .doesNotContain(ITERATION_3_COMPLETED_OUTPUT); } + + @Test + void testContainerBalancerStatusVerboseShowsFailureBreakdown() throws IOException { + ScmClient scmClient = mock(ScmClient.class); + ContainerBalancerConfiguration config = getContainerBalancerConfiguration(); + StorageContainerLocationProtocolProtos.ContainerBalancerTaskIterationStatusInfoProto iteration = + StorageContainerLocationProtocolProtos.ContainerBalancerTaskIterationStatusInfoProto.newBuilder() + .setIterationNumber(1) + .setIterationResult("ITERATION_COMPLETED") + .setIterationDuration(120L) + .setContainerMovesScheduled(5) + .setContainerMovesCompleted(3) + .setContainerMovesFailed(1) + .setContainerMovesTimeout(2) + .addContainerMoveFailures( + StorageContainerLocationProtocolProtos.ContainerMoveFailureDetailProto.newBuilder() + .setReason("REPLICATION_FAIL_TIME_OUT") + .setCount(2) + .addSourceFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid("source-uuid-1").setCount(1).build()) + .addSourceFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid("source-uuid-2").setCount(1).build()) + .addTargetFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid("target-uuid-1").setCount(1).build()) + .addTargetFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid("target-uuid-2").setCount(1).build()) + .build()) + .addContainerMoveFailures( + StorageContainerLocationProtocolProtos.ContainerMoveFailureDetailProto.newBuilder() + .setReason("PRE_MOVE_CONTAINER_NOT_FOUND") + .setCount(1) + .addSourceFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid("source-uuid-3").setCount(1).build()) + .addTargetFailureCounts( + StorageContainerLocationProtocolProtos.NodeFailureCountProto.newBuilder() + .setDatanodeUuid("target-uuid-3").setCount(1).build()) + .build()) + .build(); + + long stoppedAt = OffsetDateTime.now().toEpochSecond(); + ContainerBalancerStatusInfoProto statusInfo = ContainerBalancerStatusInfoProto.newBuilder() + .setStartedAt(stoppedAt - 600) + .setStoppedAt(stoppedAt) + .setStopReason("USER_REQUESTED") + .setStopMessage("Stopped by user request.") + .setConfiguration(config.toProtobufBuilder().setShouldRun(false)) + .addIterationsStatusInfo(iteration) + .build(); + + when(scmClient.getContainerBalancerStatusInfo()) + .thenReturn(ContainerBalancerStatusInfoResponseProto.newBuilder() + .setIsRunning(false) + .setContainerBalancerStatusInfo(statusInfo) + .build()); + + verbose.set(true); + statusCmd.execute(scmClient); + + assertThat(out.get()) + .contains("Failed container moves") + .contains("REPLICATION_FAIL_TIME_OUT") + .contains("PRE_MOVE_CONTAINER_NOT_FOUND") + .contains("source-uuid-1") + .contains("target-uuid-1") + .contains("source-uuid-3") + .contains("target-uuid-3") + .doesNotContain("Failure breakdown") + .doesNotContain("Failed move details"); + } }