Skip to content
Draft
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 @@ -551,7 +551,9 @@ private void populateContainerPathFields(KeyValueContainer container,
HddsVolume hddsVolume) throws IOException {
volumeSet.readLock();
try {
// TODO StoragePolicy Check whether need to adapt storageType
// The caller has already chosen hddsVolume for the container's storage
// type and recorded that type on the container data, so there is nothing
// storage-type specific to adapt here.
String idDir = VersionedDatanodeFeatures.ScmHA.chooseContainerPathID(
hddsVolume, clusterId);
container.populatePathFields(idDir, hddsVolume);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,17 @@
package org.apache.hadoop.ozone.container.replication;

import jakarta.annotation.Nonnull;
import jakarta.annotation.Nullable;
import java.io.IOException;
import java.io.InputStream;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.Collections;
import java.util.HashSet;
import java.util.Objects;
import java.util.Set;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.conf.StorageUnit;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
Expand Down Expand Up @@ -90,6 +93,7 @@ public boolean isAllowedContainerImport(long containerID) {
public void importContainer(long containerID, Path tarFilePath,
HddsVolume targetVolume, CopyContainerCompression compression)
throws IOException {
Objects.requireNonNull(targetVolume, "targetVolume == null");
if (!importContainerProgress.add(containerID)) {
deleteFileQuietely(tarFilePath);
String log = "Container import in progress with container Id " + containerID;
Expand All @@ -116,6 +120,10 @@ public void importContainer(long containerID, Path tarFilePath,
}
ContainerUtils.verifyContainerFileChecksum(containerData, conf);
containerData.setVolume(targetVolume);
// The descriptor carries the source volume's storage type. Record the
// type of the volume actually chosen here, so the replica reports where
// it really lives rather than where its source lived.
containerData.setStorageType(targetVolume.getStorageType());
// lastDataScanTime should be cleared for an imported container
containerData.setDataScanTimestamp(null);

Expand Down Expand Up @@ -147,13 +155,26 @@ private static void deleteFileQuietely(Path tarFilePath) {
}

HddsVolume chooseNextVolume(long spaceToReserve) throws IOException {
return chooseNextVolume(spaceToReserve, null);
}

/**
* Chooses a volume for an incoming container.
*
* @param spaceToReserve space the container needs in both tmp and dest dirs
* @param storageType storage type the container should be placed on, or
* null to allow any volume. Null is used for containers
* replicated from a datanode without storage type
* support.
*/
HddsVolume chooseNextVolume(long spaceToReserve,
@Nullable StorageType storageType) throws IOException {
// Choose volume that can hold both container in tmp and dest directory
LOG.debug("Choosing volume to reserve space : {}", spaceToReserve);
// TODO: Use the target container storage type once replication/import
// requests carry it. Null preserves the existing any-volume behavior.
LOG.debug("Choosing volume to reserve space : {}, storageType: {}",
spaceToReserve, storageType);
return volumeChoosingPolicy.chooseVolume(
StorageVolumeUtil.getHddsVolumesList(volumeSet.getVolumesList()),
spaceToReserve, null);
spaceToReserve, storageType);
}

public static Path getUntarDirectory(HddsVolume hddsVolume)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,23 @@

package org.apache.hadoop.ozone.container.replication;

import jakarta.annotation.Nullable;
import java.io.IOException;
import java.io.OutputStream;
import java.util.concurrent.CompletableFuture;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;

/**
* Client-side interface for sending a container to a target datanode.
*/
public interface ContainerUploader {
/**
* @param targetVolumeStorageType StorageType the target should place the
* container on, or null to let the target choose any volume
*/
OutputStream startUpload(long containerId, DatanodeDetails target,
CompletableFuture<Void> callback, CopyContainerCompression compression)
CompletableFuture<Void> callback, CopyContainerCompression compression,
@Nullable StorageType targetVolumeStorageType)
throws IOException;
}
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,12 @@
package org.apache.hadoop.ozone.container.replication;

import com.google.common.annotations.VisibleForTesting;
import jakarta.annotation.Nullable;
import java.io.IOException;
import java.io.OutputStream;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.DatanodeDetails.Port;
Expand Down Expand Up @@ -59,7 +61,8 @@ public GrpcContainerUploader(

@Override
public OutputStream startUpload(long containerId, DatanodeDetails target,
CompletableFuture<Void> callback, CopyContainerCompression compression) throws IOException {
CompletableFuture<Void> callback, CopyContainerCompression compression,
@Nullable StorageType targetVolumeStorageType) throws IOException {

// Get container size from local datanode instead of using passed replicateSize
Long containerSize = null;
Expand All @@ -82,7 +85,8 @@ public OutputStream startUpload(long containerId, DatanodeDetails target,
(CallStreamObserver<SendContainerRequest>) client.upload(
responseObserver), responseObserver);
return new SendContainerOutputStream(requestStream, containerId,
GrpcReplicationService.BUFFER_SIZE, compression, containerSize) {
GrpcReplicationService.BUFFER_SIZE, compression, containerSize,
targetVolumeStorageType) {
@Override
public void close() throws IOException {
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,8 @@ public void replicate(ReplicationTask task) {
CountingOutputStream output = null;
try {
output = new CountingOutputStream(
uploader.startUpload(containerID, target, fut, compression));
uploader.startUpload(containerID, target, fut, compression,
task.getTargetVolumeStorageType()));
source.copyData(containerID, output, compression);
fut.get();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@

package org.apache.hadoop.ozone.container.replication;

import jakarta.annotation.Nullable;
import java.util.Objects;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;

Expand Down Expand Up @@ -79,6 +81,16 @@ public long getContainerId() {
return cmd.getContainerID();
}

/**
* StorageType the replica should land on at the target, as chosen by SCM. Null
* when the command carries no storage type, which lets the target pick any
* volume.
*/
@Nullable
public StorageType getTargetVolumeStorageType() {
return cmd.getTargetVolumeStorageType();
}

@Override
protected Object getCommandForDebug() {
return debugString;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,9 @@

package org.apache.hadoop.ozone.container.replication;

import jakarta.annotation.Nullable;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.client.StorageTypeUtils;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerRequest;
import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
import org.apache.ratis.thirdparty.io.grpc.stub.CallStreamObserver;
Expand All @@ -28,14 +31,27 @@ class SendContainerOutputStream extends GrpcOutputStream<SendContainerRequest> {

private final CopyContainerCompression compression;
private final Long size;
/**
* StorageType the target should place the container on. Sent alongside size on
* the first request. Null lets the target choose any volume.
*/
private final StorageType targetVolumeStorageType;

SendContainerOutputStream(
CallStreamObserver<SendContainerRequest> streamObserver,
long containerId, int bufferSize, CopyContainerCompression compression,
Long size) {
this(streamObserver, containerId, bufferSize, compression, size, null);
}

SendContainerOutputStream(
CallStreamObserver<SendContainerRequest> streamObserver,
long containerId, int bufferSize, CopyContainerCompression compression,
Long size, @Nullable StorageType targetVolumeStorageType) {
super(streamObserver, containerId, bufferSize);
this.compression = compression;
this.size = size;
this.targetVolumeStorageType = targetVolumeStorageType;
}

@Override
Expand All @@ -46,9 +62,15 @@ protected void sendPart(boolean eof, int length, ByteString data) {
.setOffset(getWrittenBytes())
.setCompression(compression.toProto());

// Include container size in the first request
if (getWrittenBytes() == 0 && size != null) {
requestBuilder.setSize(size);
// Include container size and target storage type in the first request
if (getWrittenBytes() == 0) {
if (size != null) {
requestBuilder.setSize(size);
}
if (targetVolumeStorageType != null) {
requestBuilder.setStorageTypeID(
StorageTypeUtils.getID(targetVolumeStorageType));
}
}
getStreamObserver().onNext(requestBuilder.build());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@
import java.io.OutputStream;
import java.nio.file.Files;
import java.nio.file.Path;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.client.StorageTypeUtils;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerRequest;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.SendContainerResponse;
Expand Down Expand Up @@ -90,7 +92,13 @@ public void onNext(SendContainerRequest req) {
spaceToReserve = importer.getSpaceToReserve(
req.hasSize() ? req.getSize() : null);

volume = importer.chooseNextVolume(spaceToReserve);
// Keep the replica on the same storage type as its source. Absent for
// senders without storage type support, which allows any volume.
StorageType storageType = req.hasStorageTypeID()
? StorageTypeUtils.getStorageTypeFromID(req.getStorageTypeID())
: null;

volume = importer.chooseNextVolume(spaceToReserve, storageType);

Path dir = ContainerImporter.getUntarDirectory(volume);
Files.createDirectories(dir);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,10 @@

package org.apache.hadoop.ozone.protocol.commands;

import jakarta.annotation.Nullable;
import java.util.Objects;
import org.apache.hadoop.fs.StorageType;
import org.apache.hadoop.hdds.client.StorageTypeUtils;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto;
import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto.Builder;
Expand All @@ -36,6 +39,12 @@ public final class ReplicateContainerCommand
private int replicaIndex = 0;
private ReplicationCommandPriority priority =
ReplicationCommandPriority.NORMAL;
/**
* StorageType the replica should be written to on the target, so a copy stays
* on the same tier as its source. Null for containers created before storage
* policy support, which leaves the target free to choose any volume.
*/
private StorageType targetVolumeStorageType;

public static ReplicateContainerCommand toTarget(long containerID,
DatanodeDetails target) {
Expand Down Expand Up @@ -63,6 +72,10 @@ public void setPriority(ReplicationCommandPriority priority) {
this.priority = priority;
}

public void setTargetVolumeStorageType(@Nullable StorageType storageType) {
this.targetVolumeStorageType = storageType;
}

@Override
public Type getType() {
return SCMCommandProto.Type.replicateContainerCommand;
Expand All @@ -81,6 +94,10 @@ public ReplicateContainerCommandProto getProto() {
.setReplicaIndex(replicaIndex)
.setTarget(targetDatanode.getProtoBufMessage())
.setPriority(priority);
if (targetVolumeStorageType != null) {
builder.setVolumeStorageType(
StorageTypeUtils.getStorageTypeProto(targetVolumeStorageType));
}
return builder.build();
}

Expand All @@ -100,6 +117,10 @@ public static ReplicateContainerCommand getFromProtobuf(
if (protoMessage.hasPriority()) {
cmd.setPriority(protoMessage.getPriority());
}
if (protoMessage.hasVolumeStorageType()) {
cmd.setTargetVolumeStorageType(
StorageTypeUtils.getFromProtobuf(protoMessage.getVolumeStorageType()));
}
return cmd;
}

Expand All @@ -119,6 +140,11 @@ public ReplicationCommandPriority getPriority() {
return priority;
}

@Nullable
public StorageType getTargetVolumeStorageType() {
return targetVolumeStorageType;
}

@Override
public String toString() {
return getType()
Expand All @@ -129,6 +155,7 @@ public String toString() {
+ ", containerId=" + getContainerID()
+ ", replicaIndex=" + getReplicaIndex()
+ ", targetNode=" + targetDatanode
+ ", priority=" + priority;
+ ", priority=" + priority
+ ", targetVolumeStorageType=" + targetVolumeStorageType;
}
}
Loading