Author: cdouglas
Date: Tue May 3 22:50:27 2011
New Revision: 1099289
URL: http://svn.apache.org/viewvc?rev=1099289&view=rev
Log:
MAPREDUCE-2468. Add metrics for NM Shuffle.
Contributed by Luke Lu
Added:
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/mapred/
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/mapred/TestShuffleHandler.java
Modified:
hadoop/mapreduce/branches/MR-279/CHANGES.txt
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/main/java/org/apache/hadoop/mapred/ShuffleHandler.java
Modified: hadoop/mapreduce/branches/MR-279/CHANGES.txt
URL:
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/CHANGES.txt?rev=1099289&r1=1099288&r2=1099289&view=diff
==============================================================================
--- hadoop/mapreduce/branches/MR-279/CHANGES.txt (original)
+++ hadoop/mapreduce/branches/MR-279/CHANGES.txt Tue May 3 22:50:27 2011
@@ -3,6 +3,9 @@ Hadoop MapReduce Change Log
Trunk (unreleased changes)
MAPREDUCE-279
+
+ MAPREDUCE-2468. Add metrics for NM Shuffle. (Luke Lu via cdouglas)
+
Completing RM Restart. Completed Phase 3 of making sure events are logged
and restored (mahadev)
Moving userlogs out of container work-dir into a separate directory
structure. (vinodkv)
Modified:
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/main/java/org/apache/hadoop/mapred/ShuffleHandler.java
URL:
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/main/java/org/apache/hadoop/mapred/ShuffleHandler.java?rev=1099289&r1=1099288&r2=1099289&view=diff
==============================================================================
---
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/main/java/org/apache/hadoop/mapred/ShuffleHandler.java
(original)
+++
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/main/java/org/apache/hadoop/mapred/ShuffleHandler.java
Tue May 3 22:50:27 2011
@@ -22,7 +22,6 @@ import java.io.File;
import java.io.FileNotFoundException;
import java.io.IOException;
import java.io.RandomAccessFile;
-import java.io.UnsupportedEncodingException;
import java.net.URL;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
@@ -37,7 +36,6 @@ import javax.crypto.SecretKey;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.LocalDirAllocator;
-import org.apache.hadoop.mapred.IndexCache;
import org.jboss.netty.bootstrap.ServerBootstrap;
import org.jboss.netty.buffer.ChannelBuffers;
@@ -45,7 +43,6 @@ import org.jboss.netty.channel.Channel;
import org.jboss.netty.channel.ChannelFactory;
import org.jboss.netty.channel.ChannelFuture;
import org.jboss.netty.channel.ChannelFutureListener;
-import org.jboss.netty.channel.ChannelFutureProgressListener;
import org.jboss.netty.channel.ChannelHandlerContext;
import org.jboss.netty.channel.ChannelPipeline;
import org.jboss.netty.channel.ChannelPipelineFactory;
@@ -71,7 +68,6 @@ import org.jboss.netty.handler.stream.Ch
import org.jboss.netty.util.CharsetUtil;
import static org.jboss.netty.buffer.ChannelBuffers.*;
-import static org.jboss.netty.handler.codec.http.HttpHeaders.*;
import static org.jboss.netty.handler.codec.http.HttpHeaders.Names.*;
import static org.jboss.netty.handler.codec.http.HttpMethod.*;
import static org.jboss.netty.handler.codec.http.HttpResponseStatus.*;
@@ -82,35 +78,32 @@ import org.apache.commons.logging.LogFac
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.DataInputByteBuffer;
import org.apache.hadoop.io.DataOutputBuffer;
-import org.apache.hadoop.io.DataOutputByteBuffer;
import org.apache.hadoop.mapreduce.security.SecureShuffleUtils;
+import org.apache.hadoop.mapreduce.security.token.JobTokenIdentifier;
+import org.apache.hadoop.mapreduce.security.token.JobTokenSecretManager;
import org.apache.hadoop.mapreduce.task.reduce.ShuffleHeader;
-
-
+import org.apache.hadoop.metrics2.annotation.Metric;
+import org.apache.hadoop.metrics2.annotation.Metrics;
+import org.apache.hadoop.metrics2.MetricsSystem;
+import org.apache.hadoop.metrics2.lib.DefaultMetricsSystem;
+import org.apache.hadoop.metrics2.lib.MutableCounterInt;
+import org.apache.hadoop.metrics2.lib.MutableCounterLong;
+import org.apache.hadoop.metrics2.lib.MutableGaugeInt;
import org.apache.hadoop.security.token.Token;
import org.apache.hadoop.yarn.api.records.ApplicationId;
-import org.apache.hadoop.yarn.factory.providers.RecordFactoryProvider;
import org.apache.hadoop.yarn.server.nodemanager.NMConfig;
import org.apache.hadoop.yarn.server.nodemanager.containermanager.AuxServices;
import
org.apache.hadoop.yarn.server.nodemanager.containermanager.localizer.ContainerLocalizer;
import org.apache.hadoop.yarn.service.AbstractService;
import org.apache.hadoop.yarn.util.ConverterUtils;
-import org.apache.hadoop.mapreduce.security.token.JobTokenIdentifier;
-import org.apache.hadoop.mapreduce.security.token.JobTokenSecretManager;
-
-// DEBUG
-import org.apache.commons.logging.impl.Log4JLogger;
-import org.apache.log4j.Level;
+import org.apache.hadoop.yarn.util.Records;
+//
// TODO packaging
public class ShuffleHandler extends AbstractService
implements AuxServices.AuxiliaryService {
private static final Log LOG = LogFactory.getLog(ShuffleHandler.class);
- static {
- //DEBUG
- ((Log4JLogger)LOG).getLogger().setLevel(Level.DEBUG);
- }
private int port;
private ChannelFactory selector;
@@ -126,8 +119,37 @@ public class ShuffleHandler extends Abst
public static final String SHUFFLE_PORT = "mapreduce.shuffle.port";
- public ShuffleHandler() {
+ @Metrics(about="Shuffle output metrics", context="mapred")
+ static class ShuffleMetrics implements ChannelFutureListener {
+ @Metric("Shuffle output in bytes")
+ MutableCounterLong shuffleOutputBytes;
+ @Metric("# of failed shuffle outputs")
+ MutableCounterInt shuffleOutputsFailed;
+ @Metric("# of succeeeded shuffle outputs")
+ MutableCounterInt shuffleOutputsOK;
+ @Metric("# of current shuffle connections")
+ MutableGaugeInt shuffleConnections;
+
+ @Override
+ public void operationComplete(ChannelFuture future) throws Exception {
+ if (future.isSuccess()) {
+ shuffleOutputsOK.incr();
+ } else {
+ shuffleOutputsFailed.incr();
+ }
+ shuffleConnections.decr();
+ }
+ }
+
+ final ShuffleMetrics metrics;
+
+ ShuffleHandler(MetricsSystem ms) {
super("httpshuffle");
+ metrics = ms.register(new ShuffleMetrics());
+ }
+
+ public ShuffleHandler() {
+ this(DefaultMetricsSystem.instance());
}
@Override
@@ -168,7 +190,7 @@ public class ShuffleHandler extends Abst
Configuration conf = getConfig();
ServerBootstrap bootstrap = new ServerBootstrap(selector);
bootstrap.setPipelineFactory(new HttpPipelineFactory(conf));
- int port = conf.getInt("mapreduce.shuffle.port", 8080);
+ port = conf.getInt("mapreduce.shuffle.port", 8080);
accepted.add(bootstrap.bind(new InetSocketAddress(port)));
LOG.info(getName() + " listening on port " + port);
super.start();
@@ -182,14 +204,19 @@ public class ShuffleHandler extends Abst
super.stop();
}
- public static class HttpPipelineFactory implements ChannelPipelineFactory {
+ Shuffle createShuffle() {
+ return new Shuffle(getConfig());
+ }
+
+ class HttpPipelineFactory implements ChannelPipelineFactory {
- private final Shuffle SHUFFLE;
+ final Shuffle SHUFFLE;
public HttpPipelineFactory(Configuration conf) {
SHUFFLE = new Shuffle(conf);
}
+ @Override
public ChannelPipeline getPipeline() throws Exception {
return Channels.pipeline(
new HttpRequestDecoder(),
@@ -204,7 +231,7 @@ public class ShuffleHandler extends Abst
}
- static class Shuffle extends SimpleChannelUpstreamHandler {
+ class Shuffle extends SimpleChannelUpstreamHandler {
private final Configuration conf;
private final IndexCache indexCache;
@@ -216,7 +243,7 @@ public class ShuffleHandler extends Abst
indexCache = new IndexCache(new JobConf(conf));
}
- private static List<String> splitMaps(List<String> mapq) {
+ private List<String> splitMaps(List<String> mapq) {
if (null == mapq) {
return null;
}
@@ -301,6 +328,7 @@ public class ShuffleHandler extends Abst
return;
}
}
+ lastMap.addListener(metrics);
lastMap.addListener(ChannelFutureListener.CLOSE);
}
@@ -346,7 +374,7 @@ public class ShuffleHandler extends Abst
// $x/$user/appcache/$appId/output/$mapId
// TODO: Once Shuffle is out of NM, this can use MR APIs to convert
between App and Job
JobID jobID = JobID.forName(jobId);
- ApplicationId appID =
RecordFactoryProvider.getRecordFactory(null).newRecordInstance(ApplicationId.class);
+ ApplicationId appID = Records.newRecord(ApplicationId.class);
appID.setClusterTimestamp(Long.parseLong(jobID.getJtIdentifier()));
appID.setId(jobID.getId());
final String base =
@@ -388,6 +416,8 @@ public class ShuffleHandler extends Abst
partition.releaseExternalResources();
}
});
+ metrics.shuffleConnections.incr();
+ metrics.shuffleOutputBytes.incr(info.partLength); // optimistic
return writeFuture;
}
@@ -407,6 +437,7 @@ public class ShuffleHandler extends Abst
ctx.getChannel().write(response).addListener(ChannelFutureListener.CLOSE);
}
+ @Override
public void exceptionCaught(ChannelHandlerContext ctx, ExceptionEvent e)
throws Exception {
Channel ch = e.getChannel();
@@ -416,7 +447,7 @@ public class ShuffleHandler extends Abst
return;
}
- cause.printStackTrace();
+ LOG.error("Shuffle error: ", cause);
if (ch.isConnected()) {
LOG.error("Shuffle error " + e);
sendError(ctx, INTERNAL_SERVER_ERROR);
Added:
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/mapred/TestShuffleHandler.java
URL:
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/mapred/TestShuffleHandler.java?rev=1099289&view=auto
==============================================================================
---
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/mapred/TestShuffleHandler.java
(added)
+++
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-shuffle/src/test/java/org/apache/hadoop/mapred/TestShuffleHandler.java
Tue May 3 22:50:27 2011
@@ -0,0 +1,62 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.hadoop.mapred;
+
+import org.apache.hadoop.metrics2.MetricsRecordBuilder;
+import org.apache.hadoop.metrics2.MetricsSource;
+import org.apache.hadoop.metrics2.impl.MetricsSystemImpl;
+import org.apache.hadoop.metrics2.MetricsSystem;
+import static org.apache.hadoop.test.MetricsAsserts.*;
+
+import org.jboss.netty.channel.ChannelFuture;
+
+import org.junit.Test;
+import static org.apache.hadoop.test.MockitoMaker.*;
+
+public class TestShuffleHandler {
+ static final long MiB = 1024 * 1024;
+
+ @Test public void testShuffleMetrics() throws Exception {
+ MetricsSystem ms = new MetricsSystemImpl();
+ ShuffleHandler sh = new ShuffleHandler(ms);
+ ChannelFuture cf = make(stub(ChannelFuture.class).
+ returning(true, false).from.isSuccess());
+
+ sh.metrics.shuffleConnections.incr();
+ sh.metrics.shuffleOutputBytes.incr(1*MiB);
+ sh.metrics.shuffleConnections.incr();
+ sh.metrics.shuffleOutputBytes.incr(2*MiB);
+
+ checkShuffleMetrics(ms, 3*MiB, 0 , 0, 2);
+
+ sh.metrics.operationComplete(cf);
+ sh.metrics.operationComplete(cf);
+
+ checkShuffleMetrics(ms, 3*MiB, 1, 1, 0);
+ }
+
+ static void checkShuffleMetrics(MetricsSystem ms, long bytes, int failed,
+ int succeeded, int connections) {
+ MetricsSource source = ms.getSource("ShuffleMetrics");
+ MetricsRecordBuilder rb = getMetrics(source);
+ assertCounter("ShuffleOutputBytes", bytes, rb);
+ assertCounter("ShuffleOutputsFailed", failed, rb);
+ assertCounter("ShuffleOutputsOK", succeeded, rb);
+ assertGauge("ShuffleConnections", connections, rb);
+ }
+}