Rework gms beans

This commit is contained in:
elcallio 2016-10-11 14:05:47 +02:00 committed by Calle Wilund
parent e49b4ef322
commit 80762eb60a
2 changed files with 37 additions and 71 deletions

View File

@ -24,7 +24,6 @@
package org.apache.cassandra.gms; package org.apache.cassandra.gms;
import java.lang.management.ManagementFactory;
import java.net.UnknownHostException; import java.net.UnknownHostException;
import java.util.HashMap; import java.util.HashMap;
import java.util.Map; import java.util.Map;
@ -32,8 +31,6 @@ import java.util.Map;
import javax.json.JsonArray; import javax.json.JsonArray;
import javax.json.JsonObject; import javax.json.JsonObject;
import javax.json.JsonValue; import javax.json.JsonValue;
import javax.management.MBeanServer;
import javax.management.ObjectName;
import javax.management.openmbean.CompositeData; import javax.management.openmbean.CompositeData;
import javax.management.openmbean.CompositeDataSupport; import javax.management.openmbean.CompositeDataSupport;
import javax.management.openmbean.CompositeType; import javax.management.openmbean.CompositeType;
@ -45,34 +42,21 @@ import javax.management.openmbean.TabularDataSupport;
import javax.management.openmbean.TabularType; import javax.management.openmbean.TabularType;
import com.scylladb.jmx.api.APIClient; import com.scylladb.jmx.api.APIClient;
import com.scylladb.jmx.metrics.APIMBean;
public class FailureDetector implements FailureDetectorMBean { public class FailureDetector extends APIMBean implements FailureDetectorMBean {
public static final String MBEAN_NAME = "org.apache.cassandra.net:type=FailureDetector"; public static final String MBEAN_NAME = "org.apache.cassandra.net:type=FailureDetector";
private static final java.util.logging.Logger logger = java.util.logging.Logger private static final java.util.logging.Logger logger = java.util.logging.Logger
.getLogger(FailureDetector.class.getName()); .getLogger(FailureDetector.class.getName());
private APIClient c = new APIClient(); public FailureDetector(APIClient c) {
super(c);
}
public void log(String str) { public void log(String str) {
logger.finest(str); logger.finest(str);
} }
private static final FailureDetector instance = new FailureDetector();
public static FailureDetector getInstance() {
return instance;
}
private FailureDetector() {
// Register this instance with JMX
try {
MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
mbs.registerMBean(this, new ObjectName(MBEAN_NAME));
} catch (Exception e) {
throw new RuntimeException(e);
}
}
@Override @Override
public void dumpInterArrivalTimes() { public void dumpInterArrivalTimes() {
log(" dumpInterArrivalTimes()"); log(" dumpInterArrivalTimes()");
@ -86,7 +70,7 @@ public class FailureDetector implements FailureDetectorMBean {
@Override @Override
public double getPhiConvictThreshold() { public double getPhiConvictThreshold() {
log(" getPhiConvictThreshold()"); log(" getPhiConvictThreshold()");
return c.getDoubleValue("/failure_detector/phi"); return client.getDoubleValue("/failure_detector/phi");
} }
@Override @Override
@ -94,29 +78,27 @@ public class FailureDetector implements FailureDetectorMBean {
log(" getAllEndpointStates()"); log(" getAllEndpointStates()");
StringBuilder sb = new StringBuilder(); StringBuilder sb = new StringBuilder();
for (Map.Entry<String, EndpointState> entry : getEndpointStateMap().entrySet()) for (Map.Entry<String, EndpointState> entry : getEndpointStateMap().entrySet()) {
{
sb.append('/').append(entry.getKey()).append("\n"); sb.append('/').append(entry.getKey()).append("\n");
appendEndpointState(sb, entry.getValue()); appendEndpointState(sb, entry.getValue());
} }
return sb.toString(); return sb.toString();
} }
private void appendEndpointState(StringBuilder sb, EndpointState endpointState) private void appendEndpointState(StringBuilder sb, EndpointState endpointState) {
{
sb.append(" generation:").append(endpointState.getHeartBeatState().getGeneration()).append("\n"); sb.append(" generation:").append(endpointState.getHeartBeatState().getGeneration()).append("\n");
sb.append(" heartbeat:").append(endpointState.getHeartBeatState().getHeartBeatVersion()).append("\n"); sb.append(" heartbeat:").append(endpointState.getHeartBeatState().getHeartBeatVersion()).append("\n");
for (Map.Entry<ApplicationState, String> state : endpointState.applicationState.entrySet()) for (Map.Entry<ApplicationState, String> state : endpointState.applicationState.entrySet()) {
{ if (state.getKey() == ApplicationState.TOKENS) {
if (state.getKey() == ApplicationState.TOKENS)
continue; continue;
}
sb.append(" ").append(state.getKey()).append(":").append(state.getValue()).append("\n"); sb.append(" ").append(state.getKey()).append(":").append(state.getValue()).append("\n");
} }
} }
public Map<String, EndpointState> getEndpointStateMap() { public Map<String, EndpointState> getEndpointStateMap() {
Map<String, EndpointState> res = new HashMap<String, EndpointState>(); Map<String, EndpointState> res = new HashMap<String, EndpointState>();
JsonArray arr = c.getJsonArray("/failure_detector/endpoints"); JsonArray arr = client.getJsonArray("/failure_detector/endpoints");
for (int i = 0; i < arr.size(); i++) { for (int i = 0; i < arr.size(); i++) {
JsonObject obj = arr.getJsonObject(i); JsonObject obj = arr.getJsonObject(i);
EndpointState ep = new EndpointState(new HeartBeatState(obj.getInt("generation"), obj.getInt("version"))); EndpointState ep = new EndpointState(new HeartBeatState(obj.getInt("generation"), obj.getInt("version")));
@ -135,31 +117,34 @@ public class FailureDetector implements FailureDetectorMBean {
@Override @Override
public String getEndpointState(String address) throws UnknownHostException { public String getEndpointState(String address) throws UnknownHostException {
log(" getEndpointState(String address) throws UnknownHostException"); log(" getEndpointState(String address) throws UnknownHostException");
return c.getStringValue("/failure_detector/endpoints/states/" + address); return client.getStringValue("/failure_detector/endpoints/states/" + address);
} }
@Override @Override
public Map<String, String> getSimpleStates() { public Map<String, String> getSimpleStates() {
log(" getSimpleStates()"); log(" getSimpleStates()");
return c.getMapStrValue("/failure_detector/simple_states"); return client.getMapStrValue("/failure_detector/simple_states");
} }
@Override @Override
public int getDownEndpointCount() { public int getDownEndpointCount() {
log(" getDownEndpointCount()"); log(" getDownEndpointCount()");
return c.getIntValue("/failure_detector/count/endpoint/down"); return client.getIntValue("/failure_detector/count/endpoint/down");
} }
@Override @Override
public int getUpEndpointCount() { public int getUpEndpointCount() {
log(" getUpEndpointCount()"); log(" getUpEndpointCount()");
return c.getIntValue("/failure_detector/count/endpoint/up"); return client.getIntValue("/failure_detector/count/endpoint/up");
} }
// From origin: // From origin:
// this is useless except to provide backwards compatibility in phi_convict_threshold, // this is useless except to provide backwards compatibility in
// because everyone seems pretty accustomed to the default of 8, and users who have // phi_convict_threshold,
// already tuned their phi_convict_threshold for their own environments won't need to // because everyone seems pretty accustomed to the default of 8, and users
// who have
// already tuned their phi_convict_threshold for their own environments
// won't need to
// change. // change.
private final double PHI_FACTOR = 1.0 / Math.log(10.0); // 0.434... private final double PHI_FACTOR = 1.0 / Math.log(10.0); // 0.434...
@ -170,7 +155,7 @@ public class FailureDetector implements FailureDetectorMBean {
new OpenType[] { SimpleType.STRING, SimpleType.DOUBLE }); new OpenType[] { SimpleType.STRING, SimpleType.DOUBLE });
final TabularDataSupport results = new TabularDataSupport( final TabularDataSupport results = new TabularDataSupport(
new TabularType("PhiList", "PhiList", ct, new String[] { "Endpoint" })); new TabularType("PhiList", "PhiList", ct, new String[] { "Endpoint" }));
final JsonArray arr = c.getJsonArray("/failure_detector/endpoint_phi_values"); final JsonArray arr = client.getJsonArray("/failure_detector/endpoint_phi_values");
for (JsonValue v : arr) { for (JsonValue v : arr) {
JsonObject o = (JsonObject) v; JsonObject o = (JsonObject) v;
@ -187,5 +172,5 @@ public class FailureDetector implements FailureDetectorMBean {
} }
return results; return results;
} }
} }

View File

@ -23,15 +23,14 @@
*/ */
package org.apache.cassandra.gms; package org.apache.cassandra.gms;
import java.lang.management.ManagementFactory;
import java.net.UnknownHostException; import java.net.UnknownHostException;
import java.util.logging.Logger;
import javax.management.MBeanServer;
import javax.management.ObjectName;
import javax.ws.rs.core.MultivaluedHashMap; import javax.ws.rs.core.MultivaluedHashMap;
import javax.ws.rs.core.MultivaluedMap; import javax.ws.rs.core.MultivaluedMap;
import com.scylladb.jmx.api.APIClient; import com.scylladb.jmx.api.APIClient;
import com.scylladb.jmx.metrics.APIMBean;
/** /**
* This module is responsible for Gossiping information for the local endpoint. * This module is responsible for Gossiping information for the local endpoint.
@ -48,61 +47,43 @@ import com.scylladb.jmx.api.APIClient;
* node as down in the Failure Detector. * node as down in the Failure Detector.
*/ */
public class Gossiper implements GossiperMBean { public class Gossiper extends APIMBean implements GossiperMBean {
public static final String MBEAN_NAME = "org.apache.cassandra.net:type=Gossiper"; public static final String MBEAN_NAME = "org.apache.cassandra.net:type=Gossiper";
private static final java.util.logging.Logger logger = java.util.logging.Logger private static final Logger logger = Logger.getLogger(Gossiper.class.getName());
.getLogger(Gossiper.class.getName());
private APIClient c = new APIClient(); public Gossiper(APIClient c) {
super(c);
}
public void log(String str) { public void log(String str) {
logger.finest(str); logger.finest(str);
} }
private static final Gossiper instance = new Gossiper();
public static Gossiper getInstance() {
return instance;
}
private Gossiper() {
// Register this instance with JMX
try {
MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
mbs.registerMBean(this, new ObjectName(MBEAN_NAME));
} catch (Exception e) {
throw new RuntimeException(e);
}
}
@Override @Override
public long getEndpointDowntime(String address) throws UnknownHostException { public long getEndpointDowntime(String address) throws UnknownHostException {
log(" getEndpointDowntime(String address) throws UnknownHostException"); log(" getEndpointDowntime(String address) throws UnknownHostException");
return c.getLongValue("gossiper/downtime/" + address); return client.getLongValue("gossiper/downtime/" + address);
} }
@Override @Override
public int getCurrentGenerationNumber(String address) public int getCurrentGenerationNumber(String address) throws UnknownHostException {
throws UnknownHostException {
log(" getCurrentGenerationNumber(String address) throws UnknownHostException"); log(" getCurrentGenerationNumber(String address) throws UnknownHostException");
return c.getIntValue("gossiper/generation_number/" + address); return client.getIntValue("gossiper/generation_number/" + address);
} }
@Override @Override
public void unsafeAssassinateEndpoint(String address) public void unsafeAssassinateEndpoint(String address) throws UnknownHostException {
throws UnknownHostException {
log(" unsafeAssassinateEndpoint(String address) throws UnknownHostException"); log(" unsafeAssassinateEndpoint(String address) throws UnknownHostException");
MultivaluedMap<String, String> queryParams = new MultivaluedHashMap<String, String>(); MultivaluedMap<String, String> queryParams = new MultivaluedHashMap<String, String>();
queryParams.add("unsafe", "True"); queryParams.add("unsafe", "True");
c.post("gossiper/assassinate/" + address, queryParams); client.post("gossiper/assassinate/" + address, queryParams);
} }
@Override @Override
public void assassinateEndpoint(String address) throws UnknownHostException { public void assassinateEndpoint(String address) throws UnknownHostException {
log(" assassinateEndpoint(String address) throws UnknownHostException"); log(" assassinateEndpoint(String address) throws UnknownHostException");
c.post("gossiper/assassinate/" + address, null); client.post("gossiper/assassinate/" + address, null);
} }
} }