diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/HddsConfigKeys.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/HddsConfigKeys.java index 4804bfdebaca..8e6b06eb2c9d 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/HddsConfigKeys.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/HddsConfigKeys.java @@ -67,6 +67,11 @@ public final class HddsConfigKeys { "hdds.pipeline.action.max.limit"; public static final int HDDS_PIPELINE_ACTION_MAX_LIMIT_DEFAULT = 20; + // Maximum number of reports (full + incremental) a datanode includes in a single heartbeat. + public static final String HDDS_CONTAINER_REPORT_MAX_LIMIT = + "hdds.container.report.max.limit"; + public static final int HDDS_CONTAINER_REPORT_MAX_LIMIT_DEFAULT = + 10000; // Configuration to allow volume choosing policy. public static final String HDDS_DATANODE_VOLUME_CHOOSING_POLICY = "hdds.datanode.volume.choosing.policy"; diff --git a/hadoop-hdds/common/src/main/resources/ozone-default.xml b/hadoop-hdds/common/src/main/resources/ozone-default.xml index cbb095de7bfe..43874f648a06 100644 --- a/hadoop-hdds/common/src/main/resources/ozone-default.xml +++ b/hadoop-hdds/common/src/main/resources/ozone-default.xml @@ -1841,6 +1841,15 @@ single heartbeat. + + hdds.container.report.max.limit + 10000 + DATANODE + + Maximum number of reports a datanode includes in a single heartbeat to + an SCM endpoint. + + hdds.db.profile DISK diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/StateContext.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/StateContext.java index 4dcb74b40d53..78a6aef494bc 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/StateContext.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/StateContext.java @@ -375,16 +375,33 @@ public void putBackReports(List reportsToPutBack, */ public List getAllAvailableReports( HostAndPort endpoint) { - int maxLimit = Integer.MAX_VALUE; - // TODO: It is highly unlikely that we will reach maxLimit for the number - // for the number of reports, specially as it does not apply to the - // number of entries in a report. But if maxLimit is hit, should a - // heartbeat be scheduled ASAP? Should full reports not included be - // dropped? Currently this code will keep the full reports not sent - // and include it in the next heartbeat. + return getAllAvailableReportsUpToLimit(endpoint, Integer.MAX_VALUE); + } + + /** + * Returns up to {@code maxLimit} available reports for the endpoint, + * removing them from the queue. Any reports beyond the limit stay queued + * and are returned by a later call; callers can use {@link + * #hasPendingReports(HostAndPort)} to detect that case and schedule a + * follow-up heartbeat. + * + * @return List of reports + */ + public List getAllAvailableReports( + HostAndPort endpoint, int maxLimit) { return getAllAvailableReportsUpToLimit(endpoint, maxLimit); } + /** + * Returns true if incremental reports are still queued for the endpoint. + */ + public boolean hasPendingReports(HostAndPort endpoint) { + synchronized (incrementalReportsQueue) { + List reportsForEndpoint = incrementalReportsQueue.get(endpoint); + return reportsForEndpoint != null && !reportsForEndpoint.isEmpty(); + } + } + /** * Gets a point in time snapshot of all containers, any pending incremental * container reports (ICR) for containers will be included in this report diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java index 6c03d140d84f..0017f1783ab8 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java @@ -19,6 +19,8 @@ import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_ACTION_MAX_LIMIT; import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_ACTION_MAX_LIMIT_DEFAULT; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_MAX_LIMIT; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_MAX_LIMIT_DEFAULT; import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_ADDRESS_REFRESH_MISSED_COUNT_THRESHOLD; import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_ADDRESS_REFRESH_MISSED_COUNT_THRESHOLD_DEFAULT; import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_ACTION_MAX_LIMIT; @@ -81,6 +83,7 @@ public class HeartbeatEndpointTask private StateContext context; private int maxContainerActionsPerHB; private int maxPipelineActionsPerHB; + private int maxReportsPerHB; private HDDSLayoutVersionManager layoutVersionManager; private final boolean resolveOnFailureEnabled; private final int refreshThreshold; @@ -102,6 +105,8 @@ public HeartbeatEndpointTask(EndpointStateMachine rpcEndpoint, HDDS_CONTAINER_ACTION_MAX_LIMIT_DEFAULT); this.maxPipelineActionsPerHB = conf.getInt(HDDS_PIPELINE_ACTION_MAX_LIMIT, HDDS_PIPELINE_ACTION_MAX_LIMIT_DEFAULT); + this.maxReportsPerHB = conf.getInt(HDDS_CONTAINER_REPORT_MAX_LIMIT, + HDDS_CONTAINER_REPORT_MAX_LIMIT_DEFAULT); if (versionManager != null) { this.layoutVersionManager = versionManager; } else { @@ -163,6 +168,10 @@ public EndpointStateMachine.EndPointStates call() throws Exception { processResponse(response, datanodeDetailsProto); rpcEndpoint.setLastSuccessfulHeartbeat(ZonedDateTime.now()); rpcEndpoint.zeroMissedCount(); + if (context.hasPendingReports(rpcEndpoint.getAddress())) { + // addReports was capped at maxReportsPerHB and left reports queued. Trigger immediately. + context.getParent().setNextHB(Time.monotonicNow()); + } } catch (IOException ex) { Preconditions.checkState(requestBuilder != null); // put back the reports which failed to be sent @@ -218,7 +227,7 @@ private void putBackIncrementalReports( */ private void addReports(SCMHeartbeatRequestProto.Builder requestBuilder) { for (Message report : - context.getAllAvailableReports(rpcEndpoint.getAddress())) { + context.getAllAvailableReports(rpcEndpoint.getAddress(), maxReportsPerHB)) { String reportName = report.getDescriptorForType().getFullName(); for (Descriptors.FieldDescriptor descriptor : SCMHeartbeatRequestProto.getDescriptor().getFields()) { diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java index 42220c6b99fe..d7af73472e19 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java @@ -227,6 +227,27 @@ public void testReportQueueWithAddReports() throws IOException { StateContext.CONTAINER_REPORTS_PROTO_NAME); } + @Test + public void testGetAllAvailableReportsRespectsLimit() throws IOException { + StateContext ctx = createSubject(); + HostAndPort scm1 = new HostAndPort("scm1", 9001); + ctx.addEndpoint(scm1); + + // Queue 10 ICRs; no full report is refreshed, so only ICRs are returned. + batchAddIncrementalReport(ctx, + StateContext.INCREMENTAL_CONTAINER_REPORT_PROTO_NAME, 10); + + // A limited collection returns at most the limit and leaves the rest. + assertEquals(4, ctx.getAllAvailableReports(scm1, 4).size()); + assertTrue(ctx.hasPendingReports(scm1)); + assertEquals(4, ctx.getAllAvailableReports(scm1, 4).size()); + assertTrue(ctx.hasPendingReports(scm1)); + // Last chunk drains the queue. + assertEquals(2, ctx.getAllAvailableReports(scm1, 4).size()); + assertFalse(ctx.hasPendingReports(scm1)); + assertEquals(0, ctx.getAllAvailableReports(scm1, 4).size()); + } + void batchRefreshfullReports(StateContext ctx, String reportName, int count) { for (int i = 0; i < count; i++) { ctx.refreshFullReport(newMockReport(reportName)); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/endpoint/TestHeartbeatEndpointTask.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/endpoint/TestHeartbeatEndpointTask.java index 12a38d3dcfb3..d0f9ef8d35c1 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/endpoint/TestHeartbeatEndpointTask.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/endpoint/TestHeartbeatEndpointTask.java @@ -18,6 +18,7 @@ package org.apache.hadoop.ozone.container.common.states.endpoint; import static java.util.Collections.emptyList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_MAX_LIMIT; import static org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto.Type.reconcileContainerCommand; import static org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto.Type.reconstructECContainersCommand; import static org.apache.hadoop.hdds.upgrade.HDDSLayoutVersionManager.maxLayoutVersion; @@ -26,7 +27,9 @@ import static org.junit.jupiter.api.Assertions.assertNotEquals; 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.mock; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import com.google.protobuf.UnsafeByteOperations; @@ -45,6 +48,7 @@ import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.CommandStatusReportsProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerAction; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReportsProto; +import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.IncrementalContainerReportProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMHeartbeatRequestProto; @@ -395,6 +399,48 @@ public void testheartbeatWithAllReports() throws Exception { } } + @Test + public void testHeartbeatCapsReportsAndSchedulesFollowup() throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.setInt(HDDS_CONTAINER_REPORT_MAX_LIMIT, 2); + DatanodeStateMachine datanodeStateMachine = mock(DatanodeStateMachine.class); + StateContext context = new StateContext(conf, DatanodeStates.RUNNING, + datanodeStateMachine, ""); + + when(datanodeStateMachine.getQueuedCommandCount()) + .thenReturn(new EnumCounters<>(SCMCommandProto.Type.class)); + + StorageContainerDatanodeProtocolClientSideTranslatorPB scm = + mock(StorageContainerDatanodeProtocolClientSideTranslatorPB.class); + ArgumentCaptor argument = ArgumentCaptor + .forClass(SCMHeartbeatRequestProto.class); + when(scm.sendHeartbeat(argument.capture())) + .thenAnswer(invocation -> + SCMHeartbeatResponseProto.newBuilder() + .setDatanodeUUID( + ((SCMHeartbeatRequestProto) invocation.getArgument(0)) + .getDatanodeDetails().getUuid()) + .build()); + + HeartbeatEndpointTask endpointTask = getHeartbeatEndpointTask( + conf, context, scm); + context.addEndpoint(TEST_SCM_ENDPOINT); + // Queue more ICRs than the per-heartbeat limit. + for (int i = 0; i < 5; i++) { + context.addIncrementalReport( + IncrementalContainerReportProto.getDefaultInstance()); + } + + endpointTask.call(); + + // Only the capped number of ICRs is sent in this heartbeat. + SCMHeartbeatRequestProto heartbeat = argument.getValue(); + assertEquals(2, heartbeat.getIncrementalContainerReportCount()); + // The remainder stays queued and a follow-up heartbeat is scheduled now. + assertTrue(context.hasPendingReports(TEST_SCM_ENDPOINT)); + verify(datanodeStateMachine).setNextHB(anyLong()); + } + /** * Creates HeartbeatEndpointTask with the given conf, context and * StorageContainerManager client side proxy.