Updated Branches:
  refs/heads/master 3f197a101 -> 70eecd8fb

More top updates.


Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/002a0bea
Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/002a0bea
Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/002a0bea

Branch: refs/heads/master
Commit: 002a0bea510a2c60ddb9c599e648fe8fcdd70c98
Parents: 3f197a1
Author: Aaron McCurry <[email protected]>
Authored: Sun Jun 23 19:10:42 2013 -0400
Committer: Aaron McCurry <[email protected]>
Committed: Sun Jun 23 19:10:42 2013 -0400

----------------------------------------------------------------------
 .../blur/thrift/ThriftBlurControllerServer.java |   1 +
 .../blur/thrift/ThriftBlurShardServer.java      |   1 +
 .../org/apache/blur/thrift/ThriftServer.java    |  34 +++++
 .../java/org/apache/blur/shell/TopCommand.java  | 128 ++++++++++++++-----
 .../java/org/apache/blur/thrift/BlurClient.java |  37 +++++-
 .../apache/blur/thrift/BlurClientManager.java   |  11 +-
 .../apache/blur/metrics/MetricsConstants.java   |   4 +
 7 files changed, 175 insertions(+), 41 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurControllerServer.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurControllerServer.java
 
b/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurControllerServer.java
index 5a5c4a1..36bcc38 100644
--- 
a/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurControllerServer.java
+++ 
b/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurControllerServer.java
@@ -67,6 +67,7 @@ public class ThriftBlurControllerServer extends ThriftServer {
     printUlimits();
     ReporterSetup.setupReporters(configuration);
     MemoryReporter.enable();
+    setupJvmMetrics();
     ThriftServer server = createServer(serverIndex, configuration);
     server.start();
   }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurShardServer.java
----------------------------------------------------------------------
diff --git 
a/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurShardServer.java 
b/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurShardServer.java
index ff2cec6..f80cf81 100644
--- a/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurShardServer.java
+++ b/blur-core/src/main/java/org/apache/blur/thrift/ThriftBlurShardServer.java
@@ -91,6 +91,7 @@ public class ThriftBlurShardServer extends ThriftServer {
     printUlimits();
     ReporterSetup.setupReporters(configuration);
     MemoryReporter.enable();
+    setupJvmMetrics();
     ThriftServer server = createServer(serverIndex, configuration);
     server.start();
   }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-core/src/main/java/org/apache/blur/thrift/ThriftServer.java
----------------------------------------------------------------------
diff --git a/blur-core/src/main/java/org/apache/blur/thrift/ThriftServer.java 
b/blur-core/src/main/java/org/apache/blur/thrift/ThriftServer.java
index ac08285..d7e6ee4 100644
--- a/blur-core/src/main/java/org/apache/blur/thrift/ThriftServer.java
+++ b/blur-core/src/main/java/org/apache/blur/thrift/ThriftServer.java
@@ -16,10 +16,20 @@ package org.apache.blur.thrift;
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+import static org.apache.blur.metrics.MetricsConstants.HEAP_USED;
+import static org.apache.blur.metrics.MetricsConstants.JVM;
+import static org.apache.blur.metrics.MetricsConstants.LOAD_AVERAGE;
+import static org.apache.blur.metrics.MetricsConstants.ORG_APACHE_BLUR;
+import static org.apache.blur.metrics.MetricsConstants.SYSTEM;
+
 import java.io.BufferedReader;
 import java.io.IOException;
 import java.io.InputStream;
 import java.io.InputStreamReader;
+import java.lang.management.ManagementFactory;
+import java.lang.management.MemoryMXBean;
+import java.lang.management.MemoryUsage;
+import java.lang.management.OperatingSystemMXBean;
 import java.net.InetAddress;
 import java.net.InetSocketAddress;
 import java.net.UnknownHostException;
@@ -40,6 +50,10 @@ import org.apache.blur.thrift.generated.Blur;
 import org.apache.blur.thrift.generated.Blur.Iface;
 import org.apache.blur.thrift.server.TThreadedSelectorServer;
 
+import com.yammer.metrics.Metrics;
+import com.yammer.metrics.core.Gauge;
+import com.yammer.metrics.core.MetricName;
+
 public class ThriftServer {
 
   private static final Log LOG = LogFactory.getLog(ThriftServer.class);
@@ -76,6 +90,26 @@ public class ThriftServer {
     }
     reader.close();
   }
+  
+  public static void setupJvmMetrics() {
+    final MemoryMXBean memoryMXBean = ManagementFactory.getMemoryMXBean();
+    final OperatingSystemMXBean operatingSystemMXBean = 
ManagementFactory.getOperatingSystemMXBean();
+    
+    Metrics.newGauge(new MetricName(ORG_APACHE_BLUR, SYSTEM, LOAD_AVERAGE), 
new Gauge<Double>() {
+      @Override
+      public Double value() {
+        return operatingSystemMXBean.getSystemLoadAverage();
+      }
+    });
+    Metrics.newGauge(new MetricName(ORG_APACHE_BLUR, JVM, HEAP_USED), new 
Gauge<Long>() {
+      @Override
+      public Long value() {
+        MemoryUsage usage = memoryMXBean.getHeapMemoryUsage();
+        return usage.getUsed();
+      }
+    });
+  }
+
 
   public synchronized void close() {
     if (!_closed) {

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-shell/src/main/java/org/apache/blur/shell/TopCommand.java
----------------------------------------------------------------------
diff --git a/blur-shell/src/main/java/org/apache/blur/shell/TopCommand.java 
b/blur-shell/src/main/java/org/apache/blur/shell/TopCommand.java
index a5108f8..b712c07 100644
--- a/blur-shell/src/main/java/org/apache/blur/shell/TopCommand.java
+++ b/blur-shell/src/main/java/org/apache/blur/shell/TopCommand.java
@@ -26,14 +26,19 @@ import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
 import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
 
 import jline.Terminal;
 import jline.console.ConsoleReader;
 
 import org.apache.blur.thirdparty.thrift_0_9_0.TException;
-import org.apache.blur.thrift.BlurClient;
+import org.apache.blur.thrift.BlurClientManager;
+import org.apache.blur.thrift.Connection;
 import org.apache.blur.thrift.generated.Blur;
+import org.apache.blur.thrift.generated.Blur.Client;
 import org.apache.blur.thrift.generated.Blur.Iface;
 import org.apache.blur.thrift.generated.BlurException;
 import org.apache.blur.thrift.generated.Metric;
@@ -56,6 +61,9 @@ public class TopCommand extends Command {
   private static final String CH = "bc hit";
   private static final String IQ = "in qry";
   private static final String EQ = "ex qry";
+  private static final String HU = "heap usd";
+  private static final String SL = "sys load";
+
   private static final String SHARD_SERVER = "Shard Server";
   private static final String INTERNAL_QUERIES = 
"\"org.apache.blur\":type=\"Blur\",name=\"Internal Queries/s\"";
   private static final String EXTERNAL_QUERIES = 
"\"org.apache.blur\":type=\"Blur\",name=\"External Queries/s\"";
@@ -72,10 +80,14 @@ public class TopCommand extends Command {
   private static final String TABLE_COUNT = 
"\"org.apache.blur\":type=\"Blur\",scope=\"default\",name=\"Table Count\"";
   private static final String INDEX_COUNT = 
"\"org.apache.blur\":type=\"Blur\",scope=\"default\",name=\"Index Count\"";
   private static final String SEGMENT_COUNT = 
"\"org.apache.blur\":type=\"Blur\",scope=\"default\",name=\"Segment Count\"";
-  private static final double ONE_MILLION = 1000000;
-  private static final double ONE_BILLION = 1000 * ONE_MILLION;
-  private static final double ONE_TRILLION = 1000 * ONE_BILLION;
-  private static final double ONE_QUADRILLION = 1000 * ONE_TRILLION;
+  private static final String LOAD_AVERAGE = 
"\"org.apache.blur\":type=\"System\",name=\"Load Average\"";
+  private static final String HEAP_USED = 
"\"org.apache.blur\":type=\"JVM\",name=\"Heap Used\"";
+
+  private static final double ONE_THOUSAND = 1000;
+  private static final double ONE_MILLION = ONE_THOUSAND * ONE_THOUSAND;
+  private static final double ONE_BILLION = ONE_THOUSAND * ONE_MILLION;
+  private static final double ONE_TRILLION = ONE_THOUSAND * ONE_BILLION;
+  private static final double ONE_QUADRILLION = ONE_THOUSAND * ONE_TRILLION;
 
   private int _width = Integer.MAX_VALUE;
   private int _height;
@@ -97,6 +109,7 @@ public class TopCommand extends Command {
       BlurException {
 
     AtomicBoolean quit = new AtomicBoolean();
+    AtomicBoolean help = new AtomicBoolean();
 
     Map<String, String> metricNames = new HashMap<String, String>();
     metricNames.put(IQ, INTERNAL_QUERIES);
@@ -111,11 +124,14 @@ public class TopCommand extends Command {
     metricNames.put(TC, TABLE_COUNT);
     metricNames.put(IC, INDEX_COUNT);
     metricNames.put(SC, SEGMENT_COUNT);
+    metricNames.put(HU, HEAP_USED);
+    metricNames.put(SL, LOAD_AVERAGE);
 
-    Object[] labels = new Object[] { SHARD_SERVER, EQ, IQ, CH, CM, CE, CS, RO, 
RE, IM, TC, IC, SC };
+    Object[] labels = new Object[] { SHARD_SERVER, SL, HU, IM, EQ, IQ, RO, RE, 
CH, CM, CE, CS, TC, IC, SC };
 
     Set<String> sizes = new HashSet<String>();
     sizes.add(IM);
+    sizes.add(SL);
 
     Set<String> keys = new HashSet<String>(metricNames.values());
 
@@ -141,12 +157,26 @@ public class TopCommand extends Command {
           e.printStackTrace();
         }
       }
-      startCommandWatcher(reader, quit, this);
+
+      startCommandWatcher(reader, quit, help, this);
     }
 
-    Map<String, Blur.Iface> shardClients = new HashMap<String, Blur.Iface>();
+    Map<String, AtomicReference<Blur.Iface>> shardClients = new 
ConcurrentHashMap<String, AtomicReference<Blur.Iface>>();
     for (String sc : shardServerList) {
-      shardClients.put(sc, BlurClient.getClient(sc));
+      AtomicReference<Iface> ref = shardClients.get(sc);
+      if (ref == null) {
+        ref = new AtomicReference<Blur.Iface>();
+        shardClients.put(sc, ref);
+      }
+      try {
+        Client c = BlurClientManager.newClient(new Connection(sc));
+        ref.set(c);
+      } catch (IOException e) {
+        ref.set(null);
+        if (Main.debug) {
+          e.printStackTrace();
+        }
+      }
     }
 
     int longestServerName = Math.max(getSizeOfLongestKey(shardClients), 
SHARD_SERVER.length());
@@ -158,33 +188,43 @@ public class TopCommand extends Command {
     header.append("%n");
 
     do {
+      StringBuilder output = new StringBuilder();
       if (quit.get()) {
         return;
-      }
-      out.printf(truncate(header.toString()), labels);
-      for (Entry<String, Blur.Iface> e : shardClients.entrySet()) {
-        String shardServer = e.getKey();
-        Iface shardClient = e.getValue();
-        Object[] cols = new Object[labels.length];
-        int c = 0;
-        cols[c++] = shardServer;
-        StringBuilder sb = new StringBuilder("%" + longestServerName + "s");
-
-        Map<String, Metric> metrics = shardClient.metrics(keys);
-        for (int i = 1; i < labels.length; i++) {
-          String mn = metricNames.get(labels[i]);
-          Metric metric = metrics.get(mn);
-          Map<String, Double> doubleMap = metric.getDoubleMap();
-          Double value = doubleMap.get("oneMinuteRate");
-          if (value == null) {
-            value = doubleMap.get("value");
+      } else if (help.get()) {
+        showHelp(output, labels, metricNames);
+      } else {
+        output.append(truncate(String.format(header.toString(), labels)));
+        for (Entry<String, AtomicReference<Blur.Iface>> e : new 
TreeMap<String, AtomicReference<Blur.Iface>>(
+            shardClients).entrySet()) {
+          String shardServer = e.getKey();
+          Iface shardClient = e.getValue().get();
+          if (shardClient == null) {
+            String line = String.format("%" + longestServerName + "s*%n", 
shardServer);
+            output.append(line);
+          } else {
+            Object[] cols = new Object[labels.length];
+            int c = 0;
+            cols[c++] = shardServer;
+            StringBuilder sb = new StringBuilder("%" + longestServerName + 
"s");
+            Map<String, Metric> metrics = shardClient.metrics(keys);
+            for (int i = 1; i < labels.length; i++) {
+              String mn = metricNames.get(labels[i]);
+              Metric metric = metrics.get(mn);
+              Map<String, Double> doubleMap = metric.getDoubleMap();
+              Double value = doubleMap.get("oneMinuteRate");
+              if (value == null) {
+                value = doubleMap.get("value");
+              }
+              cols[c++] = humanize(value, sizes.contains(mn));
+              sb.append(" %10s");
+            }
+            sb.append("%n");
+            output.append(truncate(String.format(sb.toString(), cols)));
           }
-          cols[c++] = humanize(value, sizes.contains(mn));
-          sb.append(" %10s");
         }
-        sb.append("%n");
-        out.printf(truncate(sb.toString()), cols);
       }
+      out.print(output.toString());
       out.flush();
       if (reader != null) {
         try {
@@ -209,7 +249,22 @@ public class TopCommand extends Command {
 
   }
 
-  private void startCommandWatcher(final ConsoleReader reader, final 
AtomicBoolean quit, final Object lock) {
+  private void showHelp(StringBuilder output, Object[] labels, Map<String, 
String> metricNames) {
+    output.append("Help\n");
+
+    for (int i = 0; i < labels.length; i++) {
+      String shortName = (String) labels[i];
+      String longName = metricNames.get(shortName);
+      output.append(String.format("%15s", shortName));
+      output.append(" - ");
+      output.append(longName);
+      output.append('\n');
+    }
+
+  }
+
+  private void startCommandWatcher(final ConsoleReader reader, final 
AtomicBoolean quit, final AtomicBoolean help,
+      final Object lock) {
     Thread thread = new Thread(new Runnable() {
       @Override
       public void run() {
@@ -222,6 +277,11 @@ public class TopCommand extends Command {
                 lock.notify();
               }
               return;
+            } else if (readCharacter == 'h') {
+              help.set(!help.get());
+              synchronized (lock) {
+                lock.notify();
+              }
             }
           }
         } catch (IOException e) {
@@ -256,6 +316,10 @@ public class TopCommand extends Command {
     if (v > 0) {
       return String.format("%7.2f%s", value / ONE_MILLION, "M");
     }
+    v = (long) (value / ONE_THOUSAND);
+    if (v > 0) {
+      return String.format("%7.2f%s", value / ONE_THOUSAND, "K");
+    }
     return String.format("%7.2f", value);
   }
 

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClient.java
----------------------------------------------------------------------
diff --git a/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClient.java 
b/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClient.java
index 92a0172..b294b89 100644
--- a/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClient.java
+++ b/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClient.java
@@ -25,20 +25,31 @@ import java.util.List;
 
 import org.apache.blur.thirdparty.thrift_0_9_0.TException;
 import org.apache.blur.thrift.commands.BlurCommand;
-import org.apache.blur.thrift.generated.BlurException;
 import org.apache.blur.thrift.generated.Blur.Client;
 import org.apache.blur.thrift.generated.Blur.Iface;
+import org.apache.blur.thrift.generated.BlurException;
 
 public class BlurClient {
 
   static class BlurClientInvocationHandler implements InvocationHandler {
 
     private List<Connection> connections;
+    private int _maxRetries = BlurClientManager.MAX_RETRIES;
+    private long _backOffTime = BlurClientManager.BACK_OFF_TIME;
+    private long _maxBackOffTime = BlurClientManager.MAX_BACK_OFF_TIME;
 
     public BlurClientInvocationHandler(List<Connection> connections) {
       this.connections = connections;
     }
 
+    public BlurClientInvocationHandler(List<Connection> connections, int 
maxRetries, long backOffTime,
+        long maxBackOffTime) {
+      this(connections);
+      _maxRetries = maxRetries;
+      _backOffTime = backOffTime;
+      _maxBackOffTime = maxBackOffTime;
+    }
+
     @Override
     public Object invoke(Object proxy, final Method method, final Object[] 
args) throws Throwable {
       return BlurClientManager.execute(connections, new BlurCommand<Object>() {
@@ -61,7 +72,7 @@ public class BlurClient {
             throw new RuntimeException(targetException);
           }
         }
-      });
+      }, _maxRetries, _backOffTime, _maxBackOffTime);
     }
 
   }
@@ -73,10 +84,11 @@ public class BlurClient {
    * Blur.Iface client = 
Blur.getClient(&quot;controller1:40010,controller2:40010&quot;);
    * </pre>
    * 
-   * The connectionStr also supports passing a proxy host/port (e.g. a SOCKS 
proxy configuration):
+   * The connectionStr also supports passing a proxy host/port (e.g. a SOCKS
+   * proxy configuration):
    * 
    * <pre>
-   * Blur.Iface client = Blur.getClient("host1:port/proxyhost1:proxyport");
+   * Blur.Iface client = 
Blur.getClient(&quot;host1:port/proxyhost1:proxyport&quot;);
    * </pre>
    * 
    * @param connectionStr
@@ -88,12 +100,27 @@ public class BlurClient {
     return getClient(connections);
   }
 
+  public static Iface getClient(String connectionStr, int maxRetries, long 
backOffTime, long maxBackOffTime) {
+    List<Connection> connections = 
BlurClientManager.getConnections(connectionStr);
+    return getClient(connections, maxRetries, backOffTime, maxBackOffTime);
+  }
+
   public static Iface getClient(Connection connection) {
     return getClient(Arrays.asList(connection));
   }
 
   public static Iface getClient(List<Connection> connections) {
-    return (Iface) Proxy.newProxyInstance(Iface.class.getClassLoader(), new 
Class[] { Iface.class }, new BlurClientInvocationHandler(connections));
+    return (Iface) Proxy.newProxyInstance(Iface.class.getClassLoader(), new 
Class[] { Iface.class },
+        new BlurClientInvocationHandler(connections));
+  }
+
+  public static Iface getClient(Connection connection, int maxRetries, long 
backOffTime, long maxBackOffTime) {
+    return getClient(Arrays.asList(connection), maxRetries, backOffTime, 
maxBackOffTime);
+  }
+
+  public static Iface getClient(List<Connection> connections, int maxRetries, 
long backOffTime, long maxBackOffTime) {
+    return (Iface) Proxy.newProxyInstance(Iface.class.getClassLoader(), new 
Class[] { Iface.class },
+        new BlurClientInvocationHandler(connections, maxRetries, backOffTime, 
maxBackOffTime));
   }
 
 }

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClientManager.java
----------------------------------------------------------------------
diff --git 
a/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClientManager.java 
b/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClientManager.java
index 4e4c40c..2a25a29 100644
--- a/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClientManager.java
+++ b/blur-thrift/src/main/java/org/apache/blur/thrift/BlurClientManager.java
@@ -54,9 +54,9 @@ public class BlurClientManager {
   private static final Object NULL = new Object();
 
   private static final Log LOG = LogFactory.getLog(BlurClientManager.class);
-  private static final int MAX_RETRIES = 5;
-  private static final long BACK_OFF_TIME = 
TimeUnit.MILLISECONDS.toMillis(250);
-  private static final long MAX_BACK_OFF_TIME = TimeUnit.SECONDS.toMillis(10);
+  public static final int MAX_RETRIES = 5;
+  public static final long BACK_OFF_TIME = TimeUnit.MILLISECONDS.toMillis(250);
+  public static final long MAX_BACK_OFF_TIME = TimeUnit.SECONDS.toMillis(10);
   private static final long ONE_SECOND = TimeUnit.SECONDS.toMillis(1);
 
   private static Map<Connection, BlockingQueue<Client>> clientPool = new 
ConcurrentHashMap<Connection, BlockingQueue<Client>>();
@@ -195,7 +195,7 @@ public class BlurClientManager {
       if (allBad) {
         connectionErrorCount++;
         LOG.error("All connections are bad [" + connectionErrorCount + "].");
-        if (connectionErrorCount >= 5) {
+        if (connectionErrorCount >= maxRetries) {
           throw new IOException("All connections are bad.");
         }
         try {
@@ -236,6 +236,9 @@ public class BlurClientManager {
   }
 
   public static void sleep(long backOffTime, long maxBackOffTime, int retry, 
int maxRetries) {
+    if (maxRetries == 0) {
+      return;
+    }
     long extra = (maxBackOffTime - backOffTime) / maxRetries;
     long sleep = backOffTime + (extra * retry);
     LOG.info("Backing off call for [{0} ms]", sleep);

http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/002a0bea/blur-util/src/main/java/org/apache/blur/metrics/MetricsConstants.java
----------------------------------------------------------------------
diff --git 
a/blur-util/src/main/java/org/apache/blur/metrics/MetricsConstants.java 
b/blur-util/src/main/java/org/apache/blur/metrics/MetricsConstants.java
index 37a3fed..2e94412 100644
--- a/blur-util/src/main/java/org/apache/blur/metrics/MetricsConstants.java
+++ b/blur-util/src/main/java/org/apache/blur/metrics/MetricsConstants.java
@@ -37,6 +37,8 @@ public class MetricsConstants {
   public static final String HIT = "Hit";
   public static final String MISS = "Miss";
   public static final String CACHE = "Cache";
+  public static final String JVM = "JVM";
+  public static final String HEAP_USED = "Heap Used";
   public static final String EVICTION = "Eviction";
   public static final String TABLE_COUNT = "Table Count";
   public static final String FILES_IN_QUEUE_TO_BE_DELETED = "Files in Queue to 
be Deleted";
@@ -47,4 +49,6 @@ public class MetricsConstants {
   public static final String SEARCH_TIMER = "search-timer";
   public static final String ENTRIES = "Entries";
   public static final String SIZE = "Size";
+  public static final String LOAD_AVERAGE = "Load Average";
+  public static final String SYSTEM = "System";
 }

Reply via email to