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 @@ -36,6 +36,7 @@
import org.apache.hadoop.hdfs.protocol.datatransfer.BlockPinningException;
import org.apache.hadoop.hdfs.protocol.datatransfer.DataTransferProtoUtil;
import org.apache.hadoop.hdfs.protocol.datatransfer.IOStreamPair;
import org.apache.hadoop.hdfs.protocol.datatransfer.InvalidEncryptionKeyException;
import org.apache.hadoop.hdfs.protocol.datatransfer.Op;
import org.apache.hadoop.hdfs.protocol.datatransfer.Receiver;
import org.apache.hadoop.hdfs.protocol.datatransfer.Sender;
Expand Down Expand Up @@ -797,7 +798,6 @@ public void writeBlock(final ExtendedBlock block,
mirrorNode = targets[0].getXferAddr(connectToDnViaHostname);
LOG.debug("Connecting to datanode {}", mirrorNode);
mirrorTarget = NetUtils.createSocketAddr(mirrorNode);
mirrorSock = datanode.newSocket();
try {

DataNodeFaultInjector.get().failMirrorConnection();
Expand All @@ -806,32 +806,52 @@ public void writeBlock(final ExtendedBlock block,
(HdfsConstants.READ_TIMEOUT_EXTENSION * targets.length);
int writeTimeout = dnConf.socketWriteTimeout +
(HdfsConstants.WRITE_TIMEOUT_EXTENSION * targets.length);
NetUtils.connect(mirrorSock, mirrorTarget, timeoutValue);
mirrorSock.setTcpNoDelay(dnConf.getDataTransferServerTcpNoDelay());
mirrorSock.setSoTimeout(timeoutValue);
mirrorSock.setKeepAlive(true);
if (dnConf.getTransferSocketSendBufferSize() > 0) {
mirrorSock.setSendBufferSize(
dnConf.getTransferSocketSendBufferSize());
}

OutputStream unbufMirrorOut = NetUtils.getOutputStream(mirrorSock,
writeTimeout);
InputStream unbufMirrorIn = NetUtils.getInputStream(mirrorSock);
DataEncryptionKeyFactory keyFactory =
datanode.getDataEncryptionKeyFactoryForBlock(block);
SecretKey secretKey = null;
if (dnConf.overwriteDownstreamDerivedQOP) {
String bpid = block.getBlockPoolId();
BlockKey blockKey = datanode.blockPoolTokenSecretManager
.get(bpid).getCurrentKey();
secretKey = blockKey.getKey();
OutputStream unbufMirrorOut;
InputStream unbufMirrorIn;
int encryptionKeyRetryCount = 0;
while (true) {
try {
mirrorSock = datanode.newSocket();
NetUtils.connect(mirrorSock, mirrorTarget, timeoutValue);
mirrorSock.setTcpNoDelay(
dnConf.getDataTransferServerTcpNoDelay());
mirrorSock.setSoTimeout(timeoutValue);
mirrorSock.setKeepAlive(true);
if (dnConf.getTransferSocketSendBufferSize() > 0) {
mirrorSock.setSendBufferSize(
dnConf.getTransferSocketSendBufferSize());
}

unbufMirrorOut = NetUtils.getOutputStream(mirrorSock,
writeTimeout);
unbufMirrorIn = NetUtils.getInputStream(mirrorSock);
SecretKey secretKey = null;
if (dnConf.overwriteDownstreamDerivedQOP) {
String bpid = block.getBlockPoolId();
BlockKey blockKey = datanode.blockPoolTokenSecretManager
.get(bpid).getCurrentKey();
secretKey = blockKey.getKey();
}
IOStreamPair saslStreams = datanode.saslClient.socketSend(
mirrorSock, unbufMirrorOut, unbufMirrorIn, keyFactory,
blockToken, targets[0], secretKey);
unbufMirrorOut = saslStreams.out;
unbufMirrorIn = saslStreams.in;
break;
} catch (InvalidEncryptionKeyException e) {
IOUtils.closeSocket(mirrorSock);
mirrorSock = null;
if (!prepareRetryAfterInvalidEncryptionKey(keyFactory,
++encryptionKeyRetryCount)) {
throw e;
}
LOG.info("Retrying connection to mirror {} for block {} after "
+ "InvalidEncryptionKeyException",
targets[0], block, e);
}
}
IOStreamPair saslStreams = datanode.saslClient.socketSend(
mirrorSock, unbufMirrorOut, unbufMirrorIn, keyFactory,
blockToken, targets[0], secretKey);
unbufMirrorOut = saslStreams.out;
unbufMirrorIn = saslStreams.in;
mirrorOut = new DataOutputStream(new BufferedOutputStream(unbufMirrorOut,
smallBufferSize));
mirrorIn = new DataInputStream(unbufMirrorIn);
Expand Down Expand Up @@ -1211,21 +1231,40 @@ public void replaceBlock(final ExtendedBlock block,
final String dnAddr = proxySource.getXferAddr(connectToDnViaHostname);
LOG.debug("Connecting to datanode {}", dnAddr);
InetSocketAddress proxyAddr = NetUtils.createSocketAddr(dnAddr);
proxySock = datanode.newSocket();
NetUtils.connect(proxySock, proxyAddr, dnConf.socketTimeout);
proxySock.setTcpNoDelay(dnConf.getDataTransferServerTcpNoDelay());
proxySock.setSoTimeout(dnConf.socketTimeout);
proxySock.setKeepAlive(true);

OutputStream unbufProxyOut = NetUtils.getOutputStream(proxySock,
dnConf.socketWriteTimeout);
InputStream unbufProxyIn = NetUtils.getInputStream(proxySock);
DataEncryptionKeyFactory keyFactory =
datanode.getDataEncryptionKeyFactoryForBlock(block);
IOStreamPair saslStreams = datanode.saslClient.socketSend(proxySock,
unbufProxyOut, unbufProxyIn, keyFactory, blockToken, proxySource);
unbufProxyOut = saslStreams.out;
unbufProxyIn = saslStreams.in;
OutputStream unbufProxyOut;
InputStream unbufProxyIn;
int encryptionKeyRetryCount = 0;
while (true) {
try {
proxySock = datanode.newSocket();
NetUtils.connect(proxySock, proxyAddr, dnConf.socketTimeout);
proxySock.setTcpNoDelay(dnConf.getDataTransferServerTcpNoDelay());
proxySock.setSoTimeout(dnConf.socketTimeout);
proxySock.setKeepAlive(true);

unbufProxyOut = NetUtils.getOutputStream(proxySock,
dnConf.socketWriteTimeout);
unbufProxyIn = NetUtils.getInputStream(proxySock);
IOStreamPair saslStreams = datanode.saslClient.socketSend(
proxySock, unbufProxyOut, unbufProxyIn, keyFactory, blockToken,
proxySource);
unbufProxyOut = saslStreams.out;
unbufProxyIn = saslStreams.in;
break;
} catch (InvalidEncryptionKeyException e) {
IOUtils.closeSocket(proxySock);
proxySock = null;
if (!prepareRetryAfterInvalidEncryptionKey(keyFactory,
++encryptionKeyRetryCount)) {
throw e;
}
LOG.info("Retrying connection to proxy {} for block {} after "
+ "InvalidEncryptionKeyException",
proxySource, block, e);
}
}

proxyOut = new DataOutputStream(new BufferedOutputStream(unbufProxyOut,
smallBufferSize));
Expand Down Expand Up @@ -1313,6 +1352,15 @@ public void replaceBlock(final ExtendedBlock block,
datanode.metrics.addReplaceBlockOp(elapsed());
}

private static boolean prepareRetryAfterInvalidEncryptionKey(
DataEncryptionKeyFactory keyFactory, int retryCount) {
if (retryCount > 1) {
return false;
}
keyFactory.clearDataEncryptionKey();
return true;
}


/**
* Separated for testing.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.hadoop.hdfs.protocol.ExtendedBlock;
import org.apache.hadoop.hdfs.protocol.datatransfer.BlockConstructionStage;
import org.apache.hadoop.hdfs.protocol.datatransfer.IOStreamPair;
import org.apache.hadoop.hdfs.protocol.datatransfer.InvalidEncryptionKeyException;
import org.apache.hadoop.hdfs.protocol.datatransfer.Sender;
import org.apache.hadoop.hdfs.protocol.datatransfer.sasl.DataEncryptionKeyFactory;
import org.apache.hadoop.hdfs.security.token.block.BlockTokenIdentifier;
Expand Down Expand Up @@ -90,6 +91,15 @@ class StripedBlockWriter {
init();
}

static boolean prepareRetryAfterInvalidEncryptionKey(
DataEncryptionKeyFactory keyFactory, int retryCount) {
if (retryCount > 1) {
return false;
}
keyFactory.clearDataEncryptionKey();
return true;
}

ByteBuffer getTargetBuffer() {
return targetBuffer;
}
Expand All @@ -113,12 +123,6 @@ private void init() throws IOException {
try {
InetSocketAddress targetAddr =
stripedWriter.getSocketAddress4Transfer(target);
socket = datanode.newSocket();
NetUtils.connect(socket, targetAddr,
datanode.getDnConf().getSocketTimeout());
socket.setTcpNoDelay(
datanode.getDnConf().getDataTransferServerTcpNoDelay());
socket.setSoTimeout(datanode.getDnConf().getSocketTimeout());
DataNodeFaultInjector.get().stripedBlockWriterInit(targetBuffer);

Token<BlockTokenIdentifier> blockToken =
Expand All @@ -127,15 +131,35 @@ private void init() throws IOException {
new StorageType[]{storageType}, new String[]{storageId});

long writeTimeout = datanode.getDnConf().getSocketWriteTimeout();
OutputStream unbufOut = NetUtils.getOutputStream(socket, writeTimeout);
InputStream unbufIn = NetUtils.getInputStream(socket);
DataEncryptionKeyFactory keyFactory =
datanode.getDataEncryptionKeyFactoryForBlock(block);
IOStreamPair saslStreams = datanode.getSaslClient().socketSend(
socket, unbufOut, unbufIn, keyFactory, blockToken, target);

unbufOut = saslStreams.out;
unbufIn = saslStreams.in;
OutputStream unbufOut;
InputStream unbufIn;
int encryptionKeyRetryCount = 0;
while (true) {
try {
socket = datanode.newSocket();
NetUtils.connect(socket, targetAddr,
datanode.getDnConf().getSocketTimeout());
socket.setTcpNoDelay(
datanode.getDnConf().getDataTransferServerTcpNoDelay());
socket.setSoTimeout(datanode.getDnConf().getSocketTimeout());
unbufOut = NetUtils.getOutputStream(socket, writeTimeout);
unbufIn = NetUtils.getInputStream(socket);
IOStreamPair saslStreams = datanode.getSaslClient().socketSend(
socket, unbufOut, unbufIn, keyFactory, blockToken, target);
unbufOut = saslStreams.out;
unbufIn = saslStreams.in;
break;
} catch (InvalidEncryptionKeyException e) {
IOUtils.closeSocket(socket);
socket = null;
if (!prepareRetryAfterInvalidEncryptionKey(keyFactory,
++encryptionKeyRetryCount)) {
throw e;
}
}
}

out = new DataOutputStream(new BufferedOutputStream(unbufOut,
DFSUtilClient.getSmallBufferSize(conf)));
Expand Down
Loading
Loading