Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
9 changes: 9 additions & 0 deletions hadoop-hdds/common/src/main/resources/ozone-default.xml
Original file line number Diff line number Diff line change
Expand Up @@ -1841,6 +1841,15 @@
single heartbeat.
</description>
</property>
<property>
<name>hdds.container.report.max.limit</name>
<value>10000</value>
<tag>DATANODE</tag>
<description>
Maximum number of reports a datanode includes in a single heartbeat to
an SCM endpoint.
</description>
</property>
<property>
<name>hdds.db.profile</name>
<value>DISK</value>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -375,16 +375,33 @@ public void putBackReports(List<Message> reportsToPutBack,
*/
public List<Message> 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<Message> 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<Message> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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 {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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<SCMHeartbeatRequestProto> 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.
Expand Down