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("controller1:40010,controller2:40010"); * </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("host1:port/proxyhost1:proxyport"); * </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"; }
