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 @@ -80,6 +80,7 @@
import org.apache.ratis.statemachine.StateMachine;
import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver;
import org.apache.ratis.util.UncheckedAutoCloseable;
import org.apache.ratis.util.function.CheckedConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -807,13 +808,14 @@ private boolean isAllowed(String action) {

@Override
public StateMachine.DataChannel getStreamDataChannel(
ContainerCommandRequestProto msg)
throws StorageContainerException {
ContainerCommandRequestProto msg,
CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock)
throws StorageContainerException {
long containerID = msg.getContainerID();
Container container = getContainer(containerID);
if (container != null) {
Handler handler = getHandler(getContainerType(container));
return handler.getStreamDataChannel(container, msg);
return handler.getStreamDataChannel(container, msg, putBlock);
} else {
throw new StorageContainerException(
"ContainerID " + containerID + " does not exist",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@

package org.apache.hadoop.ozone.container.common.interfaces;

import java.io.IOException;
import java.util.Map;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
Expand All @@ -26,6 +27,7 @@
import org.apache.hadoop.ozone.container.common.transport.server.ratis.DispatcherContext;
import org.apache.ratis.statemachine.StateMachine;
import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver;
import org.apache.ratis.util.function.CheckedConsumer;

/**
* Dispatcher acts as the bridge between the transport layer and
Expand Down Expand Up @@ -87,7 +89,9 @@ void validateContainerCommand(
* When uploading using stream, get StreamDataChannel.
*/
default StateMachine.DataChannel getStreamDataChannel(
ContainerCommandRequestProto msg) throws StorageContainerException {
ContainerCommandRequestProto msg,
CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock)
throws StorageContainerException {
throw new UnsupportedOperationException(
"getStreamDataChannel not supported.");
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
import org.apache.hadoop.ozone.container.keyvalue.TarContainerPacker;
import org.apache.ratis.statemachine.StateMachine;
import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver;
import org.apache.ratis.util.function.CheckedConsumer;

/**
* Dispatcher sends ContainerCommandRequests to Handler. Each Container Type
Expand Down Expand Up @@ -93,7 +94,8 @@ public static Handler getHandlerForContainerType(
}

public abstract StateMachine.DataChannel getStreamDataChannel(
Container container, ContainerCommandRequestProto msg)
Container container, ContainerCommandRequestProto msg,
CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock)
throws StorageContainerException;

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,7 @@
import org.apache.ratis.util.JavaUtils;
import org.apache.ratis.util.LifeCycle;
import org.apache.ratis.util.TaskQueue;
import org.apache.ratis.util.function.CheckedConsumer;
import org.apache.ratis.util.function.CheckedSupplier;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -726,7 +727,17 @@ private StateMachine.DataChannel getStreamDataChannel(
requestProto.getTraceID());
}
dispatchCommand(requestProto, context); // stream init
return dispatcher.getStreamDataChannel(requestProto);
final CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock
= requestProto.getCmdType() == Type.StreamInitWithPutBlock ? this::streamPutBlock : null;
return dispatcher.getStreamDataChannel(requestProto, putBlock);
}

void streamPutBlock(ContainerCommandRequestProto request) throws IOException {
final DispatcherContext context = DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_LINK)
.setStage(DispatcherContext.WriteChunkStage.COMBINED)
.setContainer2BCSIDMap(container2BCSIDMap)
.build();
dispatchCommand(request, context);
}

@Override
Expand All @@ -741,7 +752,7 @@ public CompletableFuture<DataStream> stream(RaftClientRequest request) {
.setStage(DispatcherContext.WriteChunkStage.WRITE_DATA)
.setContainer2BCSIDMap(container2BCSIDMap)
.build();
DataChannel channel = getStreamDataChannel(requestProto, context);
final DataChannel channel = getStreamDataChannel(requestProto, context);
final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ?
getChunkExecutor(requestProto.getWriteChunk()) : null;
return new LocalStream(channel, chunkExecutor);
Expand Down Expand Up @@ -773,7 +784,10 @@ public CompletableFuture<?> link(DataStream stream, LogEntryProto entry) {

final KeyValueStreamDataChannel kvStreamDataChannel =
(KeyValueStreamDataChannel) dataChannel;
kvStreamDataChannel.setLinked();
if (!kvStreamDataChannel.link()) {
return JavaUtils.completeExceptionally(new IllegalStateException(
"PutBlock was not committed on stream close: " + kvStreamDataChannel));
}
return CompletableFuture.completedFuture(null);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,7 @@
import org.apache.ratis.statemachine.StateMachine;
import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver;
import org.apache.ratis.util.function.CheckedConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -272,7 +273,8 @@ public KeyValueHandler(ConfigurationSource config,

@Override
public StateMachine.DataChannel getStreamDataChannel(
Container container, ContainerCommandRequestProto msg)
Container container, ContainerCommandRequestProto msg,
CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock)
throws StorageContainerException {
KeyValueContainer kvContainer = (KeyValueContainer) container;
checkContainerOpen(kvContainer);
Expand All @@ -282,7 +284,7 @@ public StateMachine.DataChannel getStreamDataChannel(
BlockID.getFromProtobuf(msg.getWriteChunk().getBlockID());

return chunkManager.getStreamDataChannel(kvContainer,
blockID, metrics);
blockID, putBlock, metrics);
} else {
throw new StorageContainerException("Malformed request.",
ContainerProtos.Result.IO_EXCEPTION);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import java.util.Map;
import java.util.Objects;
import org.apache.hadoop.hdds.client.BlockID;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
import org.apache.hadoop.ozone.common.ChunkBuffer;
import org.apache.hadoop.ozone.common.ChunkBufferToByteString;
Expand All @@ -40,6 +41,7 @@
import org.apache.hadoop.ozone.container.keyvalue.interfaces.BlockManager;
import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager;
import org.apache.ratis.statemachine.StateMachine;
import org.apache.ratis.util.function.CheckedConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -79,10 +81,12 @@ public String streamInit(Container container, BlockID blockID)

@Override
public StateMachine.DataChannel getStreamDataChannel(
Container container, BlockID blockID, ContainerMetrics metrics)
Container container, BlockID blockID,
CheckedConsumer<ContainerProtos.ContainerCommandRequestProto, IOException> putBlock,
ContainerMetrics metrics)
throws StorageContainerException {
return selectHandler(container)
.getStreamDataChannel(container, blockID, metrics);
.getStreamDataChannel(container, blockID, putBlock, metrics);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@
import org.apache.hadoop.ozone.container.keyvalue.interfaces.BlockManager;
import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager;
import org.apache.ratis.statemachine.StateMachine;
import org.apache.ratis.util.function.CheckedConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -111,12 +112,13 @@ public String streamInit(Container container, BlockID blockID)

@Override
public StateMachine.DataChannel getStreamDataChannel(
Container container, BlockID blockID, ContainerMetrics metrics)
throws StorageContainerException {
Container container, BlockID blockID,
CheckedConsumer<ContainerProtos.ContainerCommandRequestProto, IOException> putBlock,
ContainerMetrics metrics) throws StorageContainerException {
checkLayoutVersion(container);
final File chunkFile = getChunkFile(container, blockID);
return new KeyValueStreamDataChannel(chunkFile,
container.getContainerData(), metrics);
container.getContainerData(), putBlock, metrics);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,19 @@
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto;
import org.apache.hadoop.hdds.ratis.ContainerCommandRequestMessage;
import org.apache.hadoop.hdds.ratis.RatisHelper;
import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
import org.apache.hadoop.hdds.scm.storage.BlockDataStreamOutput;
import org.apache.hadoop.ozone.container.common.helpers.ContainerMetrics;
import org.apache.hadoop.ozone.container.common.impl.ContainerData;
import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf;
import org.apache.ratis.util.ReferenceCountedObject;
import org.apache.ratis.util.function.CheckedConsumer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -41,19 +46,28 @@ public class KeyValueStreamDataChannel extends StreamDataChannelBase {

private final Buffers buffers = new Buffers(BlockDataStreamOutput.PUT_BLOCK_REQUEST_LENGTH_MAX);

private final AtomicReference<ContainerCommandRequestProto> putBlockRequest
= new AtomicReference<>();
private final AtomicBoolean closed = new AtomicBoolean();
private final CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock;

KeyValueStreamDataChannel(File file, ContainerData containerData,
ContainerMetrics metrics)
throws StorageContainerException {
CheckedConsumer<ContainerCommandRequestProto, IOException> putBlock,
ContainerMetrics metrics) throws StorageContainerException {
super(file, containerData, metrics);
this.putBlock = putBlock;
}

@Override
ContainerProtos.Type getType() {
return ContainerProtos.Type.StreamWrite;
}

@Override
boolean isPutBlockCommittedOnClose() {
return putBlock != null;
}

@Override
public int write(ReferenceCountedObject<ByteBuffer> referenceCounted)
throws IOException {
Expand Down Expand Up @@ -87,6 +101,10 @@ static void writeFully(ByteBuffer b, WriteMethod writeMethod)
}
}

public ContainerCommandRequestProto getPutBlockRequest() {
return putBlockRequest.get();
}

void assertOpen() throws IOException {
if (closed.get()) {
throw new IOException("Already closed: " + this);
Expand All @@ -97,7 +115,13 @@ void assertOpen() throws IOException {
public void close() throws IOException {
if (closed.compareAndSet(false, true)) {
try {
writeBuffers();
if (isPutBlockCommittedOnClose()) {
// This path requires the client appends the PutBlock at the
// end of the stream.
closeWithStreamPutBlock();
} else {
writeBuffers();
}
} finally {
super.close();
}
Expand Down Expand Up @@ -128,6 +152,21 @@ private void writeBuffers() throws IOException {
}
}

static ContainerCommandRequestProto closeBuffers(
Buffers buffers, WriteMethod writeMethod) throws IOException {
final ReferenceCountedObject<ByteBuf> ref = buffers.pollAll();
final ByteBuf buf = ref.retain();
final ContainerCommandRequestProto putBlockRequestProto;
try {
putBlockRequestProto = readPutBlockRequest(buf);
// write the remaining data
writeFully(buf.nioBuffer(), writeMethod);
} finally {
ref.release();
}
return putBlockRequestProto;
}

static int readProtoLength(ByteBuf b, int lengthIndex) {
final int readerIndex = b.readerIndex();
LOG.debug("{}, lengthIndex = {}, readerIndex = {}",
Expand Down Expand Up @@ -155,6 +194,65 @@ static void setEndIndex(ByteBuf b) {
b.writerIndex(protoIndex);
}

static ContainerCommandRequestProto readPutBlockRequest(ByteBuf b)
throws IOException {
// readerIndex protoIndex lengthIndex readerIndex+readableBytes
// V V V V
// format: |--- data ---|--- proto ---|--- proto length (4 bytes) ---|
final int readerIndex = b.readerIndex();
final int lengthIndex = readerIndex + b.readableBytes() - 4;
if (lengthIndex < readerIndex) {
throw new IOException("Buffer too short for PutBlock length: " + b.readableBytes());
}
final int protoLength = readProtoLength(b.duplicate(), lengthIndex);
final int protoIndex = lengthIndex - protoLength;
if (protoIndex < readerIndex) {
throw new IOException("Invalid PutBlock proto length: " + protoLength);
}

final ContainerCommandRequestProto proto;
try {
proto = readPutBlockRequest(b.slice(protoIndex, protoLength).nioBuffer());
} catch (Throwable t) {
RatisHelper.debug(b, "catch", LOG);
throw new IOException("Failed to readPutBlockRequest from " + b
+ ": readerIndex=" + readerIndex
+ ", protoIndex=" + protoIndex
+ ", protoLength=" + protoLength
+ ", lengthIndex=" + lengthIndex, t);
}

// set index for reading data
b.writerIndex(protoIndex);

return proto;
}

private static ContainerCommandRequestProto readPutBlockRequest(ByteBuffer b)
throws IOException {
RatisHelper.debug(b, "readPutBlockRequest", LOG);
final ByteString byteString = ByteString.copyFrom(b);

final ContainerCommandRequestProto request =
ContainerCommandRequestMessage.toProto(byteString, null);

if (!request.hasPutBlock()) {
throw new StorageContainerException(
"Malformed PutBlock request. trace ID: " + request.getTraceID(),
ContainerProtos.Result.MALFORMED_REQUEST);
}
return request;
}

private void closeWithStreamPutBlock() throws IOException {
final ContainerCommandRequestProto proto =
closeBuffers(buffers, super::writeFileChannel);
putBlockRequest.set(proto);
Preconditions.checkState(putBlock != null);
putBlock.accept(proto);
setLinked();
}

interface WriteMethod {
int applyAsInt(ByteBuffer src) throws IOException;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,8 @@ public final boolean isOpen() {
return getChannel().isOpen();
}

abstract boolean isPutBlockCommittedOnClose();

protected void assertSpaceAvailability(int requested) throws StorageContainerException {
ContainerUtils.assertSpaceAvailability(containerData.getContainerID(), containerData.getVolume(), requested);
}
Expand All @@ -102,6 +104,20 @@ public void setLinked() {
linked.set(true);
}

public boolean link() {
if (isPutBlockCommittedOnClose()) {
// The PutBlock should be commited when the steam is closed
return linked.get();
} else {
setLinked();
return true;
}
}

public boolean isLinked() {
return linked.get();
}

/**
* @return true if {@link org.apache.ratis.statemachine.StateMachine.DataChannel} is already linked.
*/
Expand Down
Loading