Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
115 changes: 109 additions & 6 deletions src/java/org/apache/cassandra/metrics/ClientMetrics.java
Original file line number Diff line number Diff line change
Expand Up @@ -18,27 +18,86 @@
*/
package org.apache.cassandra.metrics;

import static org.apache.cassandra.metrics.CassandraMetricsRegistry.Metrics;

import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.Callable;

import org.apache.cassandra.transport.Connection;
import org.apache.cassandra.transport.Server;

import com.codahale.metrics.Gauge;
import com.codahale.metrics.Meter;

import static org.apache.cassandra.metrics.CassandraMetricsRegistry.Metrics;
import com.google.common.collect.ImmutableMap;


public class ClientMetrics
{
private static final MetricNameFactory factory = new DefaultNameFactory("Client");


public static final String USER = "user";
public static final String ADDRESS = "address";
public static final String VERSION = "version";
public static final String KEYSPACE = "keyspace";
public static final String PROTOCOL = "protocol";
public static final String CIPHER = "cipher";
public static final String DRIVER_VERSION = "driverVersion";
public static final String DRIVER_NAME = "driverName";
public static final String SSL = "ssl";
public static final String REQUESTS = "requests";

public static final ClientMetrics instance = new ClientMetrics();

public boolean initialized = false;

private Collection<Server> servers;

private ClientMetrics()
{
}

public <T> void addGauge(String name, final Callable<T> provider)
public List<Connection.View> getConnectionStates()
{
if (servers == null)
return Collections.emptyList();
List<Connection.View> connections = new ArrayList<>();
for (Server s : servers)
{
connections.addAll(s.getConnectionStates());
}
return connections;
}

public int getConnectedNativeClients()
{
Metrics.register(factory.createMetricName(name), (Gauge<T>) () -> {
int ret = 0;
for (Server server : servers)
ret += server.getConnectedClients();
return ret;
}

public Map<String, Integer> getConnectedNativeClientsByUser()
{
Map<String, Integer> result = new HashMap<>();
for (Server server : servers)
{
for (Entry<String, Integer> e : server.getConnectedClientsByUser().entrySet())
{
String user = e.getKey();
result.put(user, result.getOrDefault(user, 0) + e.getValue());
}
}
return result;
}

public <T> Gauge<T> addGauge(String name, final Callable<T> provider)
{
return Metrics.register(factory.createMetricName(name), (Gauge<T>) () -> {
try
{
return provider.call();
Expand All @@ -53,4 +112,48 @@ public Meter addMeter(String name)
{
return Metrics.meter(factory.createMetricName(name));
}

public synchronized void init(Collection<Server> servers)
{
this.servers = servers;
if (initialized) return;
initialized = true;

// register metrics
addGauge("connectedNativeClients", () -> getConnectedNativeClients());
addGauge("connectedNativeClientsByUser", () -> getConnectedNativeClientsByUser());
addGauge("connections", () ->
{
List<Map<String, String>> result = new ArrayList<>();
for (Server server : servers)
{
for (Connection.View connection : server.getConnectionStates())
{
result.add(new ImmutableMap.Builder<String,String>()
.put(USER, connection.getUser())
.put(ADDRESS, connection.getAddress().toString())
.put(VERSION, String.valueOf(connection.getVersion()))
.put(REQUESTS, String.valueOf(connection.getRequests()))
.put(SSL, Boolean.toString(connection.sslEnabled()))
.put(DRIVER_NAME, connection.getDriverName().orElse("undefined"))
.put(DRIVER_VERSION, connection.getDriverVersion().orElse("undefined"))
.put(CIPHER, connection.getSSLCipher().orElse("undefined"))
.put(PROTOCOL, connection.getSSLProtocol().orElse("undefined"))
.put(KEYSPACE, connection.getKeyspace().orElse(""))
.build());
}
}
return result;
});
addGauge("clientsByProtocolVersion", () ->
{
List<Map<String, String>> result = new ArrayList<>();
for (Server server : servers)
{
result.addAll(server.getClientsByProtocolVersion());
}
return result;
});
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -115,50 +115,7 @@ synchronized void initialize()
}

// register metrics
ClientMetrics.instance.addGauge("connectedNativeClients", () ->
{
int ret = 0;
for (Server server : servers)
ret += server.getConnectedClients();
return ret;
});
ClientMetrics.instance.addGauge("connectedNativeClientsByUser", () ->
{
Map<String, Integer> result = new HashMap<>();
for (Server server : servers)
{
for (Entry<String, Integer> e : server.getConnectedClientsByUser().entrySet())
{
String user = e.getKey();
result.put(user, result.getOrDefault(user, 0) + e.getValue());
}
}
return result;
});

ClientMetrics.instance.addGauge("connections", () ->
{
List<Map<String, String>> result = new ArrayList<>();
for (Server server : servers)
{
for (Map<String, String> e : server.getConnectionStates())
{
result.add(e);
}
}
return result;
});

ClientMetrics.instance.addGauge("clientsByProtocolVersion", () ->
{
List<Map<String, String>> result = new ArrayList<>();
for (Server server : servers)
{
result.addAll(server.getClientsByProtocolVersion());
}
return result;
});

ClientMetrics.instance.init(servers);
AuthMetrics.init();

initialized = true;
Expand Down
2 changes: 1 addition & 1 deletion src/java/org/apache/cassandra/tools/NodeProbe.java
Original file line number Diff line number Diff line change
Expand Up @@ -1517,7 +1517,7 @@ public Object getCompactionMetric(String metricName)

/**
* Retrieve Proxy metrics
* @param connections, connectedNativeClients, connectedNativeClientsByUser
* @param connections, connectedNativeClients, connectedNativeClientsByUser, clientsByProtocolVersion
*/
public Object getClientMetric(String metricName)
{
Expand Down
13 changes: 11 additions & 2 deletions src/java/org/apache/cassandra/tools/nodetool/ClientStats.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import java.util.Map;
import java.util.Map.Entry;

import org.apache.cassandra.metrics.ClientMetrics;
import org.apache.cassandra.tools.NodeProbe;
import org.apache.cassandra.tools.NodeTool.NodeToolCmd;
import org.apache.cassandra.tools.nodetool.formatter.TableBuilder;
Expand Down Expand Up @@ -87,8 +88,16 @@ public void execute(NodeProbe probe)
table.add("Address", "SSL", "Cipher", "Protocol", "Version", "User", "Keyspace", "Requests", "Driver-Name", "Driver-Version");
for (Map<String, String> conn : clients)
{
table.add(conn.get("address"), conn.get("ssl"), conn.get("cipher"), conn.get("protocol"), conn.get("version"),
conn.get("user"), conn.get("keyspace"), conn.get("requests"), conn.get("driverName"), conn.get("driverVersion"));
table.add(conn.get(ClientMetrics.ADDRESS),
conn.get(ClientMetrics.SSL),
conn.get(ClientMetrics.CIPHER),
conn.get(ClientMetrics.PROTOCOL),
conn.get(ClientMetrics.VERSION),
conn.get(ClientMetrics.USER),
conn.get(ClientMetrics.KEYSPACE),
conn.get(ClientMetrics.REQUESTS),
conn.get(ClientMetrics.DRIVER_NAME),
conn.get(ClientMetrics.DRIVER_VERSION));
}
table.printTo(System.out);
System.out.println();
Expand Down
26 changes: 22 additions & 4 deletions src/java/org/apache/cassandra/transport/Connection.java
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,19 @@
*/
package org.apache.cassandra.transport;

import java.net.InetSocketAddress;
import java.util.Optional;

import io.netty.channel.Channel;
import io.netty.util.AttributeKey;

public class Connection
public abstract class Connection
{
static final AttributeKey<Connection> attributeKey = AttributeKey.valueOf("CONN");

private final Channel channel;
private final ProtocolVersion version;
private final Tracker tracker;
protected final Channel channel;
protected final ProtocolVersion version;
protected final Tracker tracker;

private volatile FrameCompressor frameCompressor;

Expand Down Expand Up @@ -64,6 +67,8 @@ public Channel channel()
return channel;
}

public abstract View view();

public interface Factory
{
Connection newConnection(Channel channel, ProtocolVersion version);
Expand All @@ -73,4 +78,17 @@ public interface Tracker
{
void addConnection(Channel ch, Connection connection);
}

public interface View {
public String getUser();
public InetSocketAddress getAddress();
public int getVersion();
public long getRequests();
public boolean sslEnabled();
public Optional<String> getDriverName();
public Optional<String> getDriverVersion();
public Optional<String> getSSLCipher();
public Optional<String> getSSLProtocol();
public Optional<String> getKeyspace();
}
}
30 changes: 8 additions & 22 deletions src/java/org/apache/cassandra/transport/Server.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import java.util.stream.Collectors;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -179,30 +181,14 @@ public Map<String, Integer> getConnectedClientsByUser()
return connectionTracker.getConnectedClientsByUser();
}

public List<Map<String, String>> getConnectionStates()
public List<Connection.View> getConnectionStates()
{
List<Map<String, String>> result = new ArrayList<>();
for(Channel c : connectionTracker.allChannels)
List<Connection.View> result = new ArrayList<>();
for (Channel c : connectionTracker.allChannels)
{
Connection connection = c.attr(Connection.attributeKey).get();
if (connection instanceof ServerConnection)
{
ServerConnection conn = (ServerConnection) connection;
SslHandler sslHandler = conn.channel().pipeline().get(SslHandler.class);

result.add(new ImmutableMap.Builder<String, String>()
.put("user", conn.getClientState().getUser().getName())
.put("keyspace", conn.getClientState().getRawKeyspace() == null ? "" : conn.getClientState().getRawKeyspace())
.put("address", conn.getClientState().getRemoteAddress().toString())
.put("version", String.valueOf(conn.getVersion().asInt()))
.put("requests", String.valueOf(conn.requests.getCount()))
.put("ssl", Boolean.toString(sslHandler == null))
.put("cipher", sslHandler != null ? sslHandler.engine().getSession().getCipherSuite() : "undefined")
.put("protocol", sslHandler != null ? sslHandler.engine().getSession().getProtocol() : "undefined")
.put("driverName", conn.getClientState().getDriverName().orElse("undefined"))
.put("driverVersion", conn.getClientState().getDriverVersion().orElse("undefined"))
.build());
}
Connection conn = c.attr(Connection.attributeKey).get();
if (conn != null)
result.add(conn.view());
}
return result;
}
Expand Down
Loading