Skip to content
Merged
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 @@ -183,7 +183,8 @@ public HeartbeatResponse sendHeartbeat(DatanodeRegistration registration,
rollingUpdateStatus = PBHelperClient.convert(resp.getRollingUpgradeStatus());
}
return new HeartbeatResponse(cmds, PBHelper.convert(resp.getHaStatus()),
rollingUpdateStatus, resp.getFullBlockReportLeaseId());
rollingUpdateStatus, resp.getFullBlockReportLeaseId(),
resp.getIsSlownode());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -853,4 +853,12 @@ boolean shouldRetryInit() {
return isAlive();
}

boolean isSlownode() {
for (BPServiceActor actor : bpServices) {
if (actor.isSlownode()) {
return true;
}
}
return false;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ enum RunningState {
private String nnId = null;
private volatile RunningState runningState = RunningState.CONNECTING;
private volatile boolean shouldServiceRun = true;
private volatile boolean isSlownode = false;
private final DataNode dn;
private final DNConf dnConf;
private long prevBlockReportId;
Expand Down Expand Up @@ -205,6 +206,7 @@ Map<String, String> getActorInfoMap() {
String.valueOf(getScheduler().getLastBlockReportTime()));
info.put("maxBlockReportSize", String.valueOf(getMaxBlockReportSize()));
info.put("maxDataLength", String.valueOf(maxDataLength));
info.put("isSlownode", String.valueOf(isSlownode));
return info;
}

Expand Down Expand Up @@ -729,6 +731,7 @@ private void offerService() throws Exception {
handleRollingUpgradeStatus(resp);
}
commandProcessingThread.enqueue(resp.getCommands());
isSlownode = resp.getIsSlownode();
}
}

Expand Down Expand Up @@ -1474,4 +1477,8 @@ void stopCommandProcessingThread() {
commandProcessingThread.interrupt();
}
}

boolean isSlownode() {
return isSlownode;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -307,4 +307,27 @@ protected BPOfferService createBPOS(
Map<String, BPOfferService> getBpByNameserviceId() {
return bpByNameserviceId;
}

boolean isSlownodeByNameserviceId(String nsId) {
if (bpByNameserviceId.containsKey(nsId)) {
return bpByNameserviceId.get(nsId).isSlownode();
}
return false;
}

boolean isSlownodeByBlockPoolId(String bpId) {
if (bpByBlockPoolId.containsKey(bpId)) {
return bpByBlockPoolId.get(bpId).isSlownode();
}
return false;
}

boolean isSlownode() {
for (BPOfferService bpOfferService : bpByBlockPoolId.values()) {
if (bpOfferService.isSlownode()) {
return true;
}
}
return false;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -3814,4 +3814,16 @@ private static boolean isWrite(BlockConstructionStage stage) {
return (stage == PIPELINE_SETUP_STREAMING_RECOVERY
|| stage == PIPELINE_SETUP_APPEND_RECOVERY);
}

boolean isSlownodeByNameserviceId(String nsId) {
return blockPoolManager.isSlownodeByNameserviceId(nsId);
}

boolean isSlownodeByBlockPoolId(String bpId) {
return blockPoolManager.isSlownodeByBlockPoolId(bpId);
}

boolean isSlownode() {
return blockPoolManager.isSlownode();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -4422,8 +4422,11 @@ nodeReg, reports, getBlockPoolId(), cacheCapacity, cacheUsed,
haContext.getState().getServiceState(),
getFSImage().getCorrectLastAppliedOrWrittenTxId());

Set<String> slownodes = DatanodeManager.getSlowNodesUuidSet();
boolean isSlownode = slownodes.contains(nodeReg.getDatanodeUuid());

return new HeartbeatResponse(cmds, haState, rollingUpgradeInfo,
blockReportLeaseId);
blockReportLeaseId, isSlownode);
} finally {
readUnlock("handleHeartbeat");
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,14 +36,23 @@ public class HeartbeatResponse {
private final RollingUpgradeStatus rollingUpdateStatus;

private final long fullBlockReportLeaseId;


private final boolean isSlownode;

public HeartbeatResponse(DatanodeCommand[] cmds,
NNHAStatusHeartbeat haStatus, RollingUpgradeStatus rollingUpdateStatus,
long fullBlockReportLeaseId) {
this(cmds, haStatus, rollingUpdateStatus, fullBlockReportLeaseId, false);
}

public HeartbeatResponse(DatanodeCommand[] cmds,
NNHAStatusHeartbeat haStatus, RollingUpgradeStatus rollingUpdateStatus,
long fullBlockReportLeaseId, boolean isSlownode) {
commands = cmds;
this.haStatus = haStatus;
this.rollingUpdateStatus = rollingUpdateStatus;
this.fullBlockReportLeaseId = fullBlockReportLeaseId;
this.isSlownode = isSlownode;
}

public DatanodeCommand[] getCommands() {
Expand All @@ -61,4 +70,8 @@ public RollingUpgradeStatus getRollingUpdateStatus() {
public long getFullBlockReportLeaseId() {
return fullBlockReportLeaseId;
}

public boolean getIsSlownode() {
return isSlownode;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,7 @@ message HeartbeatResponseProto {
optional RollingUpgradeStatusProto rollingUpgradeStatus = 3;
optional RollingUpgradeStatusProto rollingUpgradeStatusV2 = 4;
optional uint64 fullBlockReportLeaseId = 5 [ default = 0 ];
optional bool isSlownode = 6 [ default = false ];
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import static org.apache.hadoop.test.MetricsAsserts.getLongCounter;
import static org.apache.hadoop.test.MetricsAsserts.getMetrics;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertSame;
Expand Down Expand Up @@ -127,7 +128,8 @@ public class TestBPOfferService {
private final int[] heartbeatCounts = new int[3];
private DataNode mockDn;
private FsDatasetSpi<?> mockFSDataset;

private boolean isSlownode;

@Before
public void setupMocks() throws Exception {
mockNN1 = setupNNMock(0);
Expand Down Expand Up @@ -216,6 +218,23 @@ public HeartbeatResponse answer(InvocationOnMock invocation)
}
}

private class HeartbeatIsSlownodeAnswer implements Answer<HeartbeatResponse> {
private final int nnIdx;

HeartbeatIsSlownodeAnswer(int nnIdx) {
this.nnIdx = nnIdx;
}

@Override
public HeartbeatResponse answer(InvocationOnMock invocation)
throws Throwable {
HeartbeatResponse heartbeatResponse = new HeartbeatResponse(
datanodeCommands[nnIdx], mockHaStatuses[nnIdx], null,
0, isSlownode);

return heartbeatResponse;
}
}

private class HeartbeatRegisterAnswer implements Answer<HeartbeatResponse> {
private final int nnIdx;
Expand Down Expand Up @@ -1182,6 +1201,44 @@ public Object answer(InvocationOnMock invocation)
}
}

@Test(timeout = 15000)
public void testSetIsSlownode() throws Exception {
assertEquals(mockDn.isSlownode(), false);
Mockito.when(mockNN1.sendHeartbeat(
Mockito.any(DatanodeRegistration.class),
Mockito.any(StorageReport[].class),
Mockito.anyLong(),
Mockito.anyLong(),
Mockito.anyInt(),
Mockito.anyInt(),
Mockito.anyInt(),
Mockito.any(VolumeFailureSummary.class),
Mockito.anyBoolean(),
Mockito.any(SlowPeerReports.class),
Mockito.any(SlowDiskReports.class)))
.thenAnswer(new HeartbeatIsSlownodeAnswer(0));

BPOfferService bpos = setupBPOSForNNs(mockNN1);
bpos.start();

try {
waitForInitialization(bpos);

bpos.triggerHeartbeatForTests();
assertFalse(bpos.isSlownode());

isSlownode = true;
bpos.triggerHeartbeatForTests();
assertTrue(bpos.isSlownode());

isSlownode = false;
bpos.triggerHeartbeatForTests();
assertFalse(bpos.isSlownode());
} finally {
bpos.stop();
}
}

@Test(timeout = 15000)
public void testCommandProcessingThread() throws Exception {
Configuration conf = new HdfsConfiguration();
Expand Down