Author: hairong
Date: Mon Aug 9 22:42:55 2010
New Revision: 983836
URL: http://svn.apache.org/viewvc?rev=983836&view=rev
Log:
HDFS-1330. Make RPCs to DataNodes timeout. Contributed by Hairong Kuang.
Modified:
hadoop/hdfs/trunk/CHANGES.txt
hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSClient.java
hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSInputStream.java
hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/server/datanode/DataNode.java
hadoop/hdfs/trunk/src/test/hdfs/org/apache/hadoop/hdfs/server/datanode/TestInterDatanodeProtocol.java
Modified: hadoop/hdfs/trunk/CHANGES.txt
URL:
http://svn.apache.org/viewvc/hadoop/hdfs/trunk/CHANGES.txt?rev=983836&r1=983835&r2=983836&view=diff
==============================================================================
--- hadoop/hdfs/trunk/CHANGES.txt (original)
+++ hadoop/hdfs/trunk/CHANGES.txt Mon Aug 9 22:42:55 2010
@@ -27,6 +27,8 @@ Trunk (unreleased changes)
HDFS-1150. Verify datanodes' identities to clients in secure clusters.
(jghoman)
+ HDFS-1330. Make RPCs to DataNodes timeout. (hairong)
+
IMPROVEMENTS
HDFS-1096. fix for prev. commit. (boryas)
Modified: hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSClient.java
URL:
http://svn.apache.org/viewvc/hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSClient.java?rev=983836&r1=983835&r2=983836&view=diff
==============================================================================
--- hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSClient.java (original)
+++ hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSClient.java Mon Aug 9
22:42:55 2010
@@ -187,7 +187,8 @@ public class DFSClient implements FSCons
}
static ClientDatanodeProtocol createClientDatanodeProtocolProxy(
- DatanodeID datanodeid, Configuration conf, LocatedBlock locatedBlock)
+ DatanodeID datanodeid, Configuration conf, int socketTimeout,
+ LocatedBlock locatedBlock)
throws IOException {
InetSocketAddress addr = NetUtils.createSocketAddr(
datanodeid.getHost() + ":" + datanodeid.getIpcPort());
@@ -199,7 +200,7 @@ public class DFSClient implements FSCons
ticket.addToken(locatedBlock.getBlockToken());
return (ClientDatanodeProtocol)RPC.getProxy(ClientDatanodeProtocol.class,
ClientDatanodeProtocol.versionID, addr, ticket, conf, NetUtils
- .getDefaultSocketFactory(conf));
+ .getDefaultSocketFactory(conf), socketTimeout);
}
/**
Modified: hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSInputStream.java
URL:
http://svn.apache.org/viewvc/hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSInputStream.java?rev=983836&r1=983835&r2=983836&view=diff
==============================================================================
--- hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSInputStream.java
(original)
+++ hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/DFSInputStream.java Mon
Aug 9 22:42:55 2010
@@ -148,7 +148,7 @@ public class DFSInputStream extends FSIn
for(DatanodeInfo datanode : locatedblock.getLocations()) {
try {
final ClientDatanodeProtocol cdp =
DFSClient.createClientDatanodeProtocolProxy(
- datanode, dfsClient.conf, locatedblock);
+ datanode, dfsClient.conf, dfsClient.socketTimeout, locatedblock);
final long n = cdp.getReplicaVisibleLength(locatedblock.getBlock());
if (n >= 0) {
return n;
Modified:
hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/server/datanode/DataNode.java
URL:
http://svn.apache.org/viewvc/hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/server/datanode/DataNode.java?rev=983836&r1=983835&r2=983836&view=diff
==============================================================================
---
hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/server/datanode/DataNode.java
(original)
+++
hadoop/hdfs/trunk/src/java/org/apache/hadoop/hdfs/server/datanode/DataNode.java
Mon Aug 9 22:42:55 2010
@@ -530,7 +530,8 @@ public class DataNode extends Configured
}
public static InterDatanodeProtocol createInterDataNodeProtocolProxy(
- DatanodeID datanodeid, final Configuration conf) throws IOException {
+ DatanodeID datanodeid, final Configuration conf, final int socketTimeout)
+ throws IOException {
final InetSocketAddress addr = NetUtils.createSocketAddr(
datanodeid.getHost() + ":" + datanodeid.getIpcPort());
if (InterDatanodeProtocol.LOG.isDebugEnabled()) {
@@ -543,7 +544,8 @@ public class DataNode extends Configured
public InterDatanodeProtocol run() throws IOException {
return (InterDatanodeProtocol) RPC.getProxy(
InterDatanodeProtocol.class, InterDatanodeProtocol.versionID,
- addr, conf);
+ addr, UserGroupInformation.getCurrentUser(), conf,
+ NetUtils.getDefaultSocketFactory(conf), socketTimeout);
}
});
} catch (InterruptedException ie) {
@@ -1717,7 +1719,8 @@ public class DataNode extends Configured
for(DatanodeID id : datanodeids) {
try {
InterDatanodeProtocol datanode = dnRegistration.equals(id)?
- this: DataNode.createInterDataNodeProtocolProxy(id, getConf());
+ this: DataNode.createInterDataNodeProtocolProxy(id, getConf(),
+ socketTimeout);
ReplicaRecoveryInfo info = callInitReplicaRecovery(datanode, rBlock);
if (info != null &&
info.getGenerationStamp() >= block.getGenerationStamp() &&
Modified:
hadoop/hdfs/trunk/src/test/hdfs/org/apache/hadoop/hdfs/server/datanode/TestInterDatanodeProtocol.java
URL:
http://svn.apache.org/viewvc/hadoop/hdfs/trunk/src/test/hdfs/org/apache/hadoop/hdfs/server/datanode/TestInterDatanodeProtocol.java?rev=983836&r1=983835&r2=983836&view=diff
==============================================================================
---
hadoop/hdfs/trunk/src/test/hdfs/org/apache/hadoop/hdfs/server/datanode/TestInterDatanodeProtocol.java
(original)
+++
hadoop/hdfs/trunk/src/test/hdfs/org/apache/hadoop/hdfs/server/datanode/TestInterDatanodeProtocol.java
Mon Aug 9 22:42:55 2010
@@ -91,9 +91,9 @@ public class TestInterDatanodeProtocol {
assertTrue(datanodeinfo.length > 0);
//connect to a data node
- InterDatanodeProtocol idp = DataNode.createInterDataNodeProtocolProxy(
- datanodeinfo[0], conf);
DataNode datanode = cluster.getDataNode(datanodeinfo[0].getIpcPort());
+ InterDatanodeProtocol idp = DataNode.createInterDataNodeProtocolProxy(
+ datanodeinfo[0], conf, datanode.socketTimeout);
assertTrue(datanode != null);
//stop block scanner, so we could compare lastScanTime