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.