Author: suresh
Date: Mon Jul 2 23:03:17 2012
New Revision: 1356515
URL: http://svn.apache.org/viewvc?rev=1356515&view=rev
Log:
HADOOP-8533. Merging change r1356504 from trunk to branch-2
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/CHANGES.txt
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/Client.java
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/ProtobufRpcEngine.java
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RPC.java
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RpcEngine.java
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/WritableRpcEngine.java
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestIPC.java
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestRPC.java
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/CHANGES.txt
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/CHANGES.txt?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/CHANGES.txt
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/CHANGES.txt
Mon Jul 2 23:03:17 2012
@@ -60,16 +60,21 @@ Release 2.0.1-alpha - UNRELEASED
HADOOP-8463. hadoop.security.auth_to_local needs a key definition and doc.
(Madhukara Phatak via eli)
+ HADOOP-8533. Remove parallel call ununsed capability in RPC.
+ (Brandon Li via suresh)
+
BUG FIXES
HADOOP-8372. NetUtils.normalizeHostName() incorrectly handles hostname
starting with a numeric character. (Junping Du via suresh)
- HADOOP-8393. hadoop-config.sh missing variable exports, causes Yarn jobs
to fail with ClassNotFoundException MRAppMaster. (phunt via tucu)
+ HADOOP-8393. hadoop-config.sh missing variable exports, causes Yarn
+ jobs to fail with ClassNotFoundException MRAppMaster. (phunt via tucu)
HADOOP-8316. Audit logging should be disabled by default. (eli)
- HADOOP-8400. All commands warn "Kerberos krb5 configuration not found"
when security is not enabled. (tucu)
+ HADOOP-8400. All commands warn "Kerberos krb5 configuration not found"
+ when security is not enabled. (tucu)
HADOOP-8406. CompressionCodecFactory.CODEC_PROVIDERS iteration is
thread-unsafe (todd)
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/Client.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/Client.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/Client.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/Client.java
Mon Jul 2 23:03:17 2012
@@ -971,43 +971,6 @@ public class Client {
}
}
- /** Call implementation used for parallel calls. */
- private class ParallelCall extends Call {
- private ParallelResults results;
- private int index;
-
- public ParallelCall(Writable param, ParallelResults results, int index) {
- super(RPC.RpcKind.RPC_WRITABLE, param);
- this.results = results;
- this.index = index;
- }
-
- /** Deliver result to result collector. */
- protected void callComplete() {
- results.callComplete(this);
- }
- }
-
- /** Result collector for parallel calls. */
- private static class ParallelResults {
- private Writable[] values;
- private int size;
- private int count;
-
- public ParallelResults(int size) {
- this.values = new Writable[size];
- this.size = size;
- }
-
- /** Collect a result. */
- public synchronized void callComplete(ParallelCall call) {
- values[call.index] = call.getRpcResult(); // store the value
- count++; // count it
- if (count == size) // if all values are in
- notify(); // then notify waiting caller
- }
- }
-
/** Construct an IPC client whose values are of the given {@link Writable}
* class. */
public Client(Class<? extends Writable> valueClass, Configuration conf,
@@ -1209,63 +1172,6 @@ public class Client {
}
}
- /**
- * @deprecated Use {@link #call(Writable[], InetSocketAddress[],
- * Class, UserGroupInformation, Configuration)} instead
- */
- @Deprecated
- public Writable[] call(Writable[] params, InetSocketAddress[] addresses)
- throws IOException, InterruptedException {
- return call(params, addresses, null, null, conf);
- }
-
- /**
- * @deprecated Use {@link #call(Writable[], InetSocketAddress[],
- * Class, UserGroupInformation, Configuration)} instead
- */
- @Deprecated
- public Writable[] call(Writable[] params, InetSocketAddress[] addresses,
- Class<?> protocol, UserGroupInformation ticket)
- throws IOException, InterruptedException {
- return call(params, addresses, protocol, ticket, conf);
- }
-
-
- /** Makes a set of calls in parallel. Each parameter is sent to the
- * corresponding address. When all values are available, or have timed out
- * or errored, the collected results are returned in an array. The array
- * contains nulls for calls that timed out or errored. */
- public Writable[] call(Writable[] params, InetSocketAddress[] addresses,
- Class<?> protocol, UserGroupInformation ticket, Configuration conf)
- throws IOException, InterruptedException {
- if (addresses.length == 0) return new Writable[0];
-
- ParallelResults results = new ParallelResults(params.length);
- synchronized (results) {
- for (int i = 0; i < params.length; i++) {
- ParallelCall call = new ParallelCall(params[i], results, i);
- try {
- ConnectionId remoteId = ConnectionId.getConnectionId(addresses[i],
- protocol, ticket, 0, conf);
- Connection connection = getConnection(remoteId, call);
- connection.sendParam(call); // send each parameter
- } catch (IOException e) {
- // log errors
- LOG.info("Calling "+addresses[i]+" caught: " +
- e.getMessage(),e);
- results.size--; // wait for one fewer result
- }
- }
- while (results.count != results.size) {
- try {
- results.wait(); // wait for all results
- } catch (InterruptedException e) {}
- }
-
- return results.values;
- }
- }
-
// for unit testing only
@InterfaceAudience.Private
@InterfaceStability.Unstable
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/ProtobufRpcEngine.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/ProtobufRpcEngine.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/ProtobufRpcEngine.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/ProtobufRpcEngine.java
Mon Jul 2 23:03:17 2012
@@ -244,12 +244,6 @@ public class ProtobufRpcEngine implement
}
}
- @Override
- public Object[] call(Method method, Object[][] params,
- InetSocketAddress[] addrs, UserGroupInformation ticket, Configuration
conf) {
- throw new UnsupportedOperationException();
- }
-
/**
* Writable Wrapper for Protocol Buffer Requests
*/
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RPC.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RPC.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RPC.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RPC.java
Mon Jul 2 23:03:17 2012
@@ -21,7 +21,6 @@ package org.apache.hadoop.ipc;
import java.lang.reflect.Field;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Proxy;
-import java.lang.reflect.Method;
import java.net.ConnectException;
import java.net.InetSocketAddress;
@@ -627,27 +626,6 @@ public class RPC {
+ proxy.getClass());
}
- /**
- * Expert: Make multiple, parallel calls to a set of servers.
- * @deprecated Use {@link #call(Method, Object[][], InetSocketAddress[],
UserGroupInformation, Configuration)} instead
- */
- @Deprecated
- public static Object[] call(Method method, Object[][] params,
- InetSocketAddress[] addrs, Configuration conf)
- throws IOException, InterruptedException {
- return call(method, params, addrs, null, conf);
- }
-
- /** Expert: Make multiple, parallel calls to a set of servers. */
- public static Object[] call(Method method, Object[][] params,
- InetSocketAddress[] addrs,
- UserGroupInformation ticket, Configuration conf)
- throws IOException, InterruptedException {
-
- return getProtocolEngine(method.getDeclaringClass(), conf)
- .call(method, params, addrs, ticket, conf);
- }
-
/** Construct a server for a protocol implementation instance listening on a
* port and address.
* @deprecated protocol interface should be passed.
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RpcEngine.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RpcEngine.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RpcEngine.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/RpcEngine.java
Mon Jul 2 23:03:17 2012
@@ -19,7 +19,6 @@
package org.apache.hadoop.ipc;
import java.io.IOException;
-import java.lang.reflect.Method;
import java.net.InetSocketAddress;
import javax.net.SocketFactory;
@@ -44,11 +43,6 @@ public interface RpcEngine {
SocketFactory factory, int rpcTimeout,
RetryPolicy connectionRetryPolicy) throws IOException;
- /** Expert: Make multiple, parallel calls to a set of servers. */
- Object[] call(Method method, Object[][] params, InetSocketAddress[] addrs,
- UserGroupInformation ticket, Configuration conf)
- throws IOException, InterruptedException;
-
/**
* Construct a server for a protocol implementation instance.
*
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/WritableRpcEngine.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/WritableRpcEngine.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/WritableRpcEngine.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/ipc/WritableRpcEngine.java
Mon Jul 2 23:03:17 2012
@@ -20,7 +20,6 @@ package org.apache.hadoop.ipc;
import java.lang.reflect.Proxy;
import java.lang.reflect.Method;
-import java.lang.reflect.Array;
import java.lang.reflect.InvocationTargetException;
import java.net.InetSocketAddress;
@@ -274,36 +273,6 @@ public class WritableRpcEngine implement
return new ProtocolProxy<T>(protocol, proxy, true);
}
- /** Expert: Make multiple, parallel calls to a set of servers. */
- public Object[] call(Method method, Object[][] params,
- InetSocketAddress[] addrs,
- UserGroupInformation ticket, Configuration conf)
- throws IOException, InterruptedException {
-
- Invocation[] invocations = new Invocation[params.length];
- for (int i = 0; i < params.length; i++)
- invocations[i] = new Invocation(method, params[i]);
- Client client = CLIENTS.getClient(conf);
- try {
- Writable[] wrappedValues =
- client.call(invocations, addrs, method.getDeclaringClass(), ticket,
conf);
-
- if (method.getReturnType() == Void.TYPE) {
- return null;
- }
-
- Object[] values =
- (Object[])Array.newInstance(method.getReturnType(),
wrappedValues.length);
- for (int i = 0; i < values.length; i++)
- if (wrappedValues[i] != null)
- values[i] = ((ObjectWritable)wrappedValues[i]).get();
-
- return values;
- } finally {
- CLIENTS.stopClient(client);
- }
- }
-
/* Construct a server for a protocol implementation instance listening on a
* port and address. */
@Override
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestIPC.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestIPC.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestIPC.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestIPC.java
Mon Jul 2 23:03:17 2012
@@ -149,41 +149,6 @@ public class TestIPC {
}
}
- private static class ParallelCaller extends Thread {
- private Client client;
- private int count;
- private InetSocketAddress[] addresses;
- private boolean failed;
-
- public ParallelCaller(Client client, InetSocketAddress[] addresses,
- int count) {
- this.client = client;
- this.addresses = addresses;
- this.count = count;
- }
-
- public void run() {
- for (int i = 0; i < count; i++) {
- try {
- Writable[] params = new Writable[addresses.length];
- for (int j = 0; j < addresses.length; j++)
- params[j] = new LongWritable(RANDOM.nextLong());
- Writable[] values = client.call(params, addresses, null, null, conf);
- for (int j = 0; j < addresses.length; j++) {
- if (!params[j].equals(values[j])) {
- LOG.fatal("Call failed!");
- failed = true;
- break;
- }
- }
- } catch (Exception e) {
- LOG.fatal("Caught: " + StringUtils.stringifyException(e));
- failed = true;
- }
- }
- }
- }
-
@Test
public void testSerial() throws Exception {
testSerial(3, false, 2, 5, 100);
@@ -218,51 +183,7 @@ public class TestIPC {
}
@Test
- public void testParallel() throws Exception {
- testParallel(10, false, 2, 4, 2, 4, 100);
- }
-
- public void testParallel(int handlerCount, boolean handlerSleep,
- int serverCount, int addressCount,
- int clientCount, int callerCount, int callCount)
- throws Exception {
- Server[] servers = new Server[serverCount];
- for (int i = 0; i < serverCount; i++) {
- servers[i] = new TestServer(handlerCount, handlerSleep);
- servers[i].start();
- }
-
- InetSocketAddress[] addresses = new InetSocketAddress[addressCount];
- for (int i = 0; i < addressCount; i++) {
- addresses[i] = NetUtils.getConnectAddress(servers[i%serverCount]);
- }
-
- Client[] clients = new Client[clientCount];
- for (int i = 0; i < clientCount; i++) {
- clients[i] = new Client(LongWritable.class, conf);
- }
-
- ParallelCaller[] callers = new ParallelCaller[callerCount];
- for (int i = 0; i < callerCount; i++) {
- callers[i] =
- new ParallelCaller(clients[i%clientCount], addresses, callCount);
- callers[i].start();
- }
- for (int i = 0; i < callerCount; i++) {
- callers[i].join();
- assertFalse(callers[i].failed);
- }
- for (int i = 0; i < clientCount; i++) {
- clients[i].stop();
- }
- for (int i = 0; i < serverCount; i++) {
- servers[i].stop();
- }
- }
-
- @Test
public void testStandAloneClient() throws Exception {
- testParallel(10, false, 2, 4, 2, 4, 100);
Client client = new Client(LongWritable.class, conf);
InetSocketAddress address = new InetSocketAddress("127.0.0.1", 10);
try {
@@ -781,13 +702,4 @@ public class TestIPC {
Ints.toByteArray(HADOOP0_21_ERROR_MSG.length()),
HADOOP0_21_ERROR_MSG.getBytes());
}
-
- public static void main(String[] args) throws Exception {
-
- //new TestIPC().testSerial(5, false, 2, 10, 1000);
-
- new TestIPC().testParallel(10, false, 2, 4, 2, 4, 1000);
-
- }
-
}
Modified:
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestRPC.java
URL:
http://svn.apache.org/viewvc/hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestRPC.java?rev=1356515&r1=1356514&r2=1356515&view=diff
==============================================================================
---
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestRPC.java
(original)
+++
hadoop/common/branches/branch-2/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/ipc/TestRPC.java
Mon Jul 2 23:03:17 2012
@@ -244,13 +244,6 @@ public class TestRPC {
*/
private static class StoppedRpcEngine implements RpcEngine {
- @Override
- public Object[] call(Method method, Object[][] params, InetSocketAddress[]
addrs,
- UserGroupInformation ticket, Configuration conf)
- throws IOException, InterruptedException {
- return null;
- }
-
@SuppressWarnings("unchecked")
@Override
public <T> ProtocolProxy<T> getProxy(Class<T> protocol, long clientVersion,
@@ -491,17 +484,6 @@ public class TestRPC {
}
}
- // try some multi-calls
- Method echo =
- TestProtocol.class.getMethod("echo", new Class[] { String.class });
- String[] strings = (String[])RPC.call(echo, new String[][]{{"a"},{"b"}},
- new InetSocketAddress[] {addr,
addr}, conf);
- assertTrue(Arrays.equals(strings, new String[]{"a","b"}));
-
- Method ping = TestProtocol.class.getMethod("ping", new Class[] {});
- Object[] voids = RPC.call(ping, new Object[][]{{},{}},
- new InetSocketAddress[] {addr, addr}, conf);
- assertEquals(voids, null);
} finally {
server.stop();
if(proxy!=null) RPC.stopProxy(proxy);