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 @@ -18,7 +18,7 @@
package org.apache.hadoop.ozone.container.common.statemachine.commandhandler;

import java.util.concurrent.atomic.AtomicLong;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.FinalizeNewLayoutVersionCommandProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.FinalizeNewDatanodeVersionCommandProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
import org.apache.hadoop.metrics2.lib.MetricsRegistry;
import org.apache.hadoop.metrics2.lib.MutableRate;
Expand Down Expand Up @@ -51,7 +51,7 @@ public FinalizeVersionCommandHandler() {
MetricsRegistry registry = new MetricsRegistry(
FinalizeVersionCommandHandler.class.getSimpleName());
this.opsLatencyMs =
registry.newRate(SCMCommandProto.Type.finalizeNewLayoutVersionCommand + "Ms");
registry.newRate(SCMCommandProto.Type.finalizeNewDatanodeVersionCommand + "Ms");
}

/**
Expand All @@ -69,10 +69,10 @@ public void handle(SCMCommand<?> command, OzoneContainer ozoneContainer,
invocationCount.incrementAndGet();
final long startTime = Time.monotonicNow();
DatanodeStateMachine dsm = context.getParent();
final FinalizeNewLayoutVersionCommandProto finalizeCommand =
final FinalizeNewDatanodeVersionCommandProto finalizeCommand =
((FinalizeVersionCommand) command).getProto();
try {
if (finalizeCommand.getFinalizeNewLayoutVersion()) {
if (finalizeCommand.getFinalizeNewDatanodeVersion()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should probably add a check that DN and SCM software versions match here, the open question is what to do about it. I'm currently thinking we should crash the node, since there's already safeguards checking versions when the DN registers and before SCM finalizes itself. With those in place it shouldn't even be possible to hit this code block so I think we can be defensive.

if (dsm.getVersionManager().needsFinalization()) {
LOG.info("Finalize upgrade called.");
dsm.getVersionManager().finalizeUpgrade();
Expand All @@ -93,7 +93,7 @@ public void handle(SCMCommand<?> command, OzoneContainer ozoneContainer,
*/
@Override
public SCMCommandProto.Type getCommandType() {
return SCMCommandProto.Type.finalizeNewLayoutVersionCommand;
return SCMCommandProto.Type.finalizeNewDatanodeVersionCommand;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.CommandQueueReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerAction;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerActionsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineAction;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineActionsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
Expand Down Expand Up @@ -128,13 +128,13 @@ public EndpointStateMachine.EndPointStates call() throws Exception {
try {
Preconditions.checkState(this.datanodeDetailsProto != null);

LayoutVersionProto versionInfo = toVersionProto(
DatanodeVersionProto versionInfo = toVersionProto(
versionManager.getApparentVersion(),
versionManager.getSoftwareVersion());

requestBuilder = SCMHeartbeatRequestProto.newBuilder()
.setDatanodeDetails(datanodeDetailsProto)
.setDataNodeLayoutVersion(versionInfo);
.setDatanodeVersion(versionInfo);
addReports(requestBuilder);
addContainerActions(requestBuilder);
addPipelineActions(requestBuilder);
Expand Down Expand Up @@ -363,9 +363,9 @@ private void processResponse(SCMHeartbeatResponseProto response,
processCommonCommand(commandResponseProto,
setNodeOperationalStateCommand);
break;
case finalizeNewLayoutVersionCommand:
case finalizeNewDatanodeVersionCommand:
FinalizeVersionCommand finalizeVersionCommand =
FinalizeVersionCommand.getFromProtobuf(commandResponseProto.getFinalizeNewLayoutVersionCommandProto());
FinalizeVersionCommand.getFromProtobuf(commandResponseProto.getFinalizeNewDatanodeVersionCommandProto());
if (LOG.isDebugEnabled()) {
LOG.debug("Received SCM finalize command {}", finalizeVersionCommand.getId());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMRegisteredResponseProto;
Expand Down Expand Up @@ -108,10 +108,10 @@ public EndpointStateMachine.EndPointStates call() throws Exception {

if (rpcEndPoint.getState()
.equals(EndpointStateMachine.EndPointStates.REGISTER)) {
LayoutVersionProto layoutInfo = LayoutVersionProto.newBuilder()
.setMetadataLayoutVersion(
DatanodeVersionProto versionInfo = DatanodeVersionProto.newBuilder()
.setApparentVersion(
versionManager.getApparentVersion().serialize())
.setSoftwareLayoutVersion(
.setSoftwareVersion(
versionManager.getSoftwareVersion().serialize())
.build();
ContainerReportsProto containerReport =
Expand All @@ -122,7 +122,7 @@ public EndpointStateMachine.EndPointStates call() throws Exception {
// TODO : Add responses to the command Queue.
SCMRegisteredResponseProto response = rpcEndPoint.getEndPoint()
.register(datanodeDetails.getExtendedProtoBufMessage(),
nodeReport, containerReport, pipelineReportsProto, layoutInfo);
nodeReport, containerReport, pipelineReportsProto, versionInfo);
Preconditions.assertEquals(datanodeDetails.getUuidString(), response.getDatanodeUUID(), "datanodeID");
Preconditions.assertTrue(!StringUtils.isBlank(response.getClusterID()),
"Invalid cluster ID in the response.");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@

import org.apache.hadoop.hdds.ComponentVersion;
import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;

/**
* Util methods for upgrade.
Expand All @@ -29,17 +29,18 @@ public final class UpgradeUtils {
private UpgradeUtils() {
}

public static LayoutVersionProto defaultVersionProto() {
public static DatanodeVersionProto defaultVersionProto() {
int softwareVersion = HDDSVersion.SOFTWARE_VERSION.serialize();
return LayoutVersionProto.newBuilder()
.setMetadataLayoutVersion(softwareVersion)
.setSoftwareLayoutVersion(softwareVersion).build();
return DatanodeVersionProto.newBuilder()
.setApparentVersion(softwareVersion)
.setSoftwareVersion(softwareVersion).build();
}

public static LayoutVersionProto toVersionProto(ComponentVersion apparentVersion, ComponentVersion softwareVersion) {
return LayoutVersionProto.newBuilder()
.setMetadataLayoutVersion(apparentVersion.serialize())
.setSoftwareLayoutVersion(softwareVersion.serialize())
public static DatanodeVersionProto toVersionProto(ComponentVersion apparentVersion,
ComponentVersion softwareVersion) {
return DatanodeVersionProto.newBuilder()
.setApparentVersion(apparentVersion.serialize())
.setSoftwareVersion(softwareVersion.serialize())
.build();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
import org.apache.hadoop.hdds.annotation.InterfaceAudience;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ExtendedDatanodeDetailsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMHeartbeatRequestProto;
Expand Down Expand Up @@ -69,14 +69,14 @@ SCMHeartbeatResponseProto sendHeartbeat(SCMHeartbeatRequestProto heartbeat)
* @param extendedDatanodeDetailsProto - extended Datanode Details.
* @param nodeReport - Node Report.
* @param containerReportsRequestProto - Container Reports.
* @param layoutInfo - Layout Version Information.
* @param versionInfo - Datanode Version Information.
* @return SCM Command.
*/
SCMRegisteredResponseProto register(
ExtendedDatanodeDetailsProto extendedDatanodeDetailsProto,
NodeReportProto nodeReport,
ContainerReportsProto containerReportsRequestProto,
PipelineReportsProto pipelineReports,
LayoutVersionProto layoutInfo) throws IOException;
DatanodeVersionProto versionInfo) throws IOException;

}
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
import org.apache.hadoop.hdds.annotation.InterfaceAudience;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.CommandQueueReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMVersionRequestProto;
Expand Down Expand Up @@ -52,13 +52,13 @@ public interface StorageContainerNodeProtocol {
* @param datanodeDetails DatanodeDetails
* @param nodeReport NodeReportProto
* @param pipelineReport PipelineReportsProto
* @param layoutVersionInfo LayoutVersionProto
* @param versionInfo DatanodeVersionProto
* @return SCMRegisteredResponseProto
*/
RegisteredCommand register(DatanodeDetails datanodeDetails,
NodeReportProto nodeReport,
PipelineReportsProto pipelineReport,
LayoutVersionProto layoutVersionInfo);
DatanodeVersionProto versionInfo);

/**
* Send heartbeat to indicate the datanode is alive and doing well.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,31 +18,31 @@
package org.apache.hadoop.ozone.protocol.commands;

import java.util.Objects;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.FinalizeNewLayoutVersionCommandProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.FinalizeNewDatanodeVersionCommandProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;

/**
* Asks DataNode to finalize new upgrade version.
*/
public class FinalizeVersionCommand
extends SCMCommand<FinalizeNewLayoutVersionCommandProto> {
extends SCMCommand<FinalizeNewDatanodeVersionCommandProto> {

private boolean finalizeUpgrade = false;
private LayoutVersionProto versionInfo;
private DatanodeVersionProto versionInfo;

public FinalizeVersionCommand(boolean finalizeNewLayoutVersion,
LayoutVersionProto versionInfo,
public FinalizeVersionCommand(boolean finalizeNewDatanodeVersion,
DatanodeVersionProto versionInfo,
long id) {
super(id);
finalizeUpgrade = finalizeNewLayoutVersion;
finalizeUpgrade = finalizeNewDatanodeVersion;
this.versionInfo = versionInfo;
}

public FinalizeVersionCommand(boolean finalizeNewLayoutVersion,
LayoutVersionProto versionInfo) {
public FinalizeVersionCommand(boolean finalizeNewDatanodeVersion,
DatanodeVersionProto versionInfo) {
super();
finalizeUpgrade = finalizeNewLayoutVersion;
finalizeUpgrade = finalizeNewDatanodeVersion;
this.versionInfo = versionInfo;
}

Expand All @@ -53,24 +53,24 @@ public FinalizeVersionCommand(boolean finalizeNewLayoutVersion,
*/
@Override
public SCMCommandProto.Type getType() {
return SCMCommandProto.Type.finalizeNewLayoutVersionCommand;
return SCMCommandProto.Type.finalizeNewDatanodeVersionCommand;
}

@Override
public FinalizeNewLayoutVersionCommandProto getProto() {
return FinalizeNewLayoutVersionCommandProto.newBuilder()
.setFinalizeNewLayoutVersion(finalizeUpgrade)
public FinalizeNewDatanodeVersionCommandProto getProto() {
return FinalizeNewDatanodeVersionCommandProto.newBuilder()
.setFinalizeNewDatanodeVersion(finalizeUpgrade)
.setCmdId(getId())
.setDataNodeLayoutVersion(versionInfo)
.setDatanodeVersion(versionInfo)
.build();
}

public static FinalizeVersionCommand getFromProtobuf(
FinalizeNewLayoutVersionCommandProto finalizeProto) {
FinalizeNewDatanodeVersionCommandProto finalizeProto) {
Objects.requireNonNull(finalizeProto, "finalizeProto == null");
return new FinalizeVersionCommand(
finalizeProto.getFinalizeNewLayoutVersion(),
finalizeProto.getDataNodeLayoutVersion(), finalizeProto.getCmdId());
finalizeProto.getFinalizeNewDatanodeVersion(),
finalizeProto.getDatanodeVersion(), finalizeProto.getCmdId());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@
import java.util.function.Consumer;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ExtendedDatanodeDetailsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMDatanodeRequest;
Expand Down Expand Up @@ -144,7 +144,7 @@ public SCMHeartbeatResponseProto sendHeartbeat(
* @param extendedDatanodeDetailsProto - extended Datanode Details
* @param nodeReport - Node Report.
* @param containerReportsRequestProto - Container Reports.
* @param layoutInfo - Layout Version Information.
* @param versionInfo - Datanode Version Information.
* @return SCM Command.
*/
@Override
Expand All @@ -153,16 +153,16 @@ public SCMRegisteredResponseProto register(
NodeReportProto nodeReport,
ContainerReportsProto containerReportsRequestProto,
PipelineReportsProto pipelineReportsProto,
LayoutVersionProto layoutInfo)
DatanodeVersionProto versionInfo)
throws IOException {
SCMRegisterRequestProto.Builder req =
SCMRegisterRequestProto.newBuilder();
req.setExtendedDatanodeDetails(extendedDatanodeDetailsProto);
req.setContainerReport(containerReportsRequestProto);
req.setPipelineReports(pipelineReportsProto);
req.setNodeReport(nodeReport);
if (layoutInfo != null) {
req.setDataNodeLayoutVersion(layoutInfo);
if (versionInfo != null) {
req.setDatanodeVersion(versionInfo);
}
return submitRequest(Type.Register,
(builder) -> builder.setRegisterRequest(req))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
import java.io.IOException;
import java.util.concurrent.TimeoutException;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMDatanodeRequest;
Expand Down Expand Up @@ -71,9 +71,9 @@ public SCMRegisteredResponseProto register(
.getContainerReport();
NodeReportProto dnNodeReport = request.getNodeReport();
PipelineReportsProto pipelineReport = request.getPipelineReports();
LayoutVersionProto versionInfo = null;
if (request.hasDataNodeLayoutVersion()) {
versionInfo = request.getDataNodeLayoutVersion();
DatanodeVersionProto versionInfo = null;
if (request.hasDatanodeVersion()) {
versionInfo = request.getDatanodeVersion();
} else {
// Backward compatibility to make sure old Datanodes can still talk to
// SCM.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.CommandStatusReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
Expand Down Expand Up @@ -227,7 +227,7 @@ private void sleepIfNeeded() {
NodeReportProto nodeReport,
ContainerReportsProto containerReportsRequestProto,
PipelineReportsProto pipelineReportsProto,
LayoutVersionProto layoutInfo)
DatanodeVersionProto versionInfo)
throws IOException {
rpcCount.incrementAndGet();
DatanodeDetailsProto datanodeDetailsProto =
Expand Down
Loading