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);
+  }
+}


Reply via email to