This is an automated email from the ASF dual-hosted git repository.

dcapwell pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 876ac8c  When a table attempts to clean up metrics, it was cleaning up 
all global table metrics
876ac8c is described below

commit 876ac8c6111697c8105899889bf794438110bc74
Author: David Capwell <[email protected]>
AuthorDate: Wed Dec 2 13:35:56 2020 -0800

    When a table attempts to clean up metrics, it was cleaning up all global 
table metrics
    
    patch by David Capwell; reviewed by Jon Meredith, Jordan West, Yifan Cai 
for CASSANDRA-16095
---
 CHANGES.txt                                        |   1 +
 .../config/CassandraRelevantProperties.java        |  71 +++++-
 .../org/apache/cassandra/db/ColumnFamilyStore.java |  22 +-
 .../org/apache/cassandra/metrics/TableMetrics.java |  78 +++---
 .../org/apache/cassandra/utils/MBeanWrapper.java   |  95 +++----
 .../distributed/impl/AbstractCluster.java          |   8 +
 .../cassandra/distributed/impl/Instance.java       |   9 +-
 .../distributed/test/metric/TableMetricTest.java   | 272 +++++++++++++++++++++
 8 files changed, 463 insertions(+), 93 deletions(-)

diff --git a/CHANGES.txt b/CHANGES.txt
index 7458801..b539e2b 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -15,6 +15,7 @@
  * Throw BufferOverflowException from DataOutputBuffer for better visibility 
(CASSANDRA-16214)
  * TLS connections to the storage port on a node without server encryption 
configured causes java.io.IOException accessing missing keystore 
(CASSANDRA-16144)
  * Internode messaging catches OOMs and does not rethrow (CASSANDRA-15214)
+ * When a table attempts to clean up metrics, it was cleaning up all global 
table metrics (CASSANDRA-16095)
 Merged from 3.11:
  * SASI's `max_compaction_flush_memory_in_mb` settings over 100GB revert to 
default of 1GB (CASSANDRA-16071)
 Merged from 3.0:
diff --git 
a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java 
b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java
index 7402aa1..5d918a8 100644
--- a/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java
+++ b/src/java/org/apache/cassandra/config/CassandraRelevantProperties.java
@@ -163,7 +163,10 @@ public enum CassandraRelevantProperties
     
ORG_APACHE_CASSANDRA_DB_VIRTUAL_SYSTEM_PROPERTIES_TABLE_TEST("org.apache.cassandra.db.virtual.SystemPropertiesTableTest"),
 
     /** This property indicates whether disable_mbean_registration is true */
-    
IS_DISABLED_MBEAN_REGISTRATION("org.apache.cassandra.disable_mbean_registration");
+    
IS_DISABLED_MBEAN_REGISTRATION("org.apache.cassandra.disable_mbean_registration"),
+
+    /** what class to use for mbean registeration */
+    MBEAN_REGISTRATION_CLASS("org.apache.cassandra.mbean_registration_class");
 
     CassandraRelevantProperties(String key, String defaultVal)
     {
@@ -197,6 +200,28 @@ public enum CassandraRelevantProperties
     }
 
     /**
+     * Gets the value of a system property as a String.
+     * @return system property String value if it exists, overrideDefaultValue 
otherwise.
+     */
+    public String getString(String overrideDefaultValue)
+    {
+        String value = System.getProperty(key);
+        if (value == null)
+            return overrideDefaultValue;
+
+        return STRING_CONVERTER.convert(value);
+    }
+
+    /**
+     * Sets the value into system properties.
+     * @param value to set
+     */
+    public void setString(String value)
+    {
+        System.setProperty(key, value);
+    }
+
+    /**
      * Gets the value of a system property as a boolean.
      * @return system property boolean value if it exists, false otherwise().
      */
@@ -208,6 +233,28 @@ public enum CassandraRelevantProperties
     }
 
     /**
+     * Gets the value of a system property as a boolean.
+     * @return system property boolean value if it exists, 
overrideDefaultValue otherwise.
+     */
+    public boolean getBoolean(boolean overrideDefaultValue)
+    {
+        String value = System.getProperty(key);
+        if (value == null)
+            return overrideDefaultValue;
+
+        return BOOLEAN_CONVERTER.convert(value);
+    }
+
+    /**
+     * Sets the value into system properties.
+     * @param value to set
+     */
+    public void setBoolean(boolean value)
+    {
+        System.setProperty(key, Boolean.toString(value));
+    }
+
+    /**
      * Gets the value of a system property as a int.
      * @return system property int value if it exists, defaultValue otherwise.
      */
@@ -218,6 +265,28 @@ public enum CassandraRelevantProperties
         return INTEGER_CONVERTER.convert(value == null ? defaultVal : value);
     }
 
+    /**
+     * Gets the value of a system property as a int.
+     * @return system property int value if it exists, overrideDefaultValue 
otherwise.
+     */
+    public int getInt(int overrideDefaultValue)
+    {
+        String value = System.getProperty(key);
+        if (value == null)
+            return overrideDefaultValue;
+
+        return INTEGER_CONVERTER.convert(value);
+    }
+
+    /**
+     * Sets the value into system properties.
+     * @param value to set
+     */
+    public void setInt(int value)
+    {
+        System.setProperty(key, Integer.toString(value));
+    }
+
     private interface PropertyConverter<T>
     {
         T convert(String value);
diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java 
b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
index f712d1b..ceacd0d 100644
--- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
+++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
@@ -439,12 +439,8 @@ public class ColumnFamilyStore implements 
ColumnFamilyStoreMBean
         if (registerBookeeping)
         {
             // register the mbean
-            mbeanName = 
String.format("org.apache.cassandra.db:type=%s,keyspace=%s,table=%s",
-                                         isIndex() ? "IndexTables" : "Tables",
-                                         keyspace.getName(), name);
-            oldMBeanName = 
String.format("org.apache.cassandra.db:type=%s,keyspace=%s,columnfamily=%s",
-                                         isIndex() ? "IndexColumnFamilies" : 
"ColumnFamilies",
-                                         keyspace.getName(), name);
+            mbeanName = getTableMBeanName(keyspace.getName(), name, isIndex());
+            oldMBeanName = getColumnFamilieMBeanName(keyspace.getName(), name, 
isIndex());
 
             String[] objectNames = {mbeanName, oldMBeanName};
             for (String objectName : objectNames)
@@ -461,6 +457,20 @@ public class ColumnFamilyStore implements 
ColumnFamilyStoreMBean
         sstableImporter = new SSTableImporter(this);
     }
 
+    public static String getTableMBeanName(String ks, String name, boolean 
isIndex)
+    {
+        return 
String.format("org.apache.cassandra.db:type=%s,keyspace=%s,table=%s",
+                      isIndex ? "IndexTables" : "Tables",
+                      ks, name);
+    }
+
+    public static String getColumnFamilieMBeanName(String ks, String name, 
boolean isIndex)
+    {
+       return 
String.format("org.apache.cassandra.db:type=%s,keyspace=%s,columnfamily=%s",
+                            isIndex ? "IndexColumnFamilies" : "ColumnFamilies",
+                            ks, name);
+    }
+
     public void updateSpeculationThreshold()
     {
         try
diff --git a/src/java/org/apache/cassandra/metrics/TableMetrics.java 
b/src/java/org/apache/cassandra/metrics/TableMetrics.java
index eb42331..634d929 100644
--- a/src/java/org/apache/cassandra/metrics/TableMetrics.java
+++ b/src/java/org/apache/cassandra/metrics/TableMetrics.java
@@ -62,6 +62,18 @@ import com.codahale.metrics.RatioGauge;
  */
 public class TableMetrics
 {
+    /**
+     * stores metrics that will be rolled into a single global metric
+     */
+    private static final ConcurrentMap<String, Set<Metric>> ALL_TABLE_METRICS 
= Maps.newConcurrentMap();
+    public static final long[] EMPTY = new long[0];
+    private static final MetricNameFactory GLOBAL_FACTORY = new 
AllTableMetricNameFactory("Table");
+    private static final MetricNameFactory GLOBAL_ALIAS_FACTORY = new 
AllTableMetricNameFactory("ColumnFamily");
+
+    public final static LatencyMetrics GLOBAL_READ_LATENCY = new 
LatencyMetrics(GLOBAL_FACTORY, GLOBAL_ALIAS_FACTORY, "Read");
+    public final static LatencyMetrics GLOBAL_WRITE_LATENCY = new 
LatencyMetrics(GLOBAL_FACTORY, GLOBAL_ALIAS_FACTORY, "Write");
+    public final static LatencyMetrics GLOBAL_RANGE_LATENCY = new 
LatencyMetrics(GLOBAL_FACTORY, GLOBAL_ALIAS_FACTORY, "Range");
+
     /** Total amount of data stored in the memtable that resides on-heap, 
including column related overhead and partitions overwritten. */
     public final Gauge<Long> memtableOnHeapDataSize;
     /** Total amount of data stored in the memtable that resides off-heap, 
including column related overhead and partitions overwritten. */
@@ -205,8 +217,6 @@ public class TableMetrics
 
     private final MetricNameFactory factory;
     private final MetricNameFactory aliasFactory;
-    private static final MetricNameFactory globalFactory = new 
AllTableMetricNameFactory("Table");
-    private static final MetricNameFactory globalAliasFactory = new 
AllTableMetricNameFactory("ColumnFamily");
 
     public final Counter speculativeRetries;
     public final Counter speculativeFailedRetries;
@@ -238,10 +248,6 @@ public class TableMetrics
     public final TableHistogram repairedDataTrackingOverreadRows;
     public final TableTimer repairedDataTrackingOverreadTime;
 
-    public final static LatencyMetrics globalReadLatency = new 
LatencyMetrics(globalFactory, globalAliasFactory, "Read");
-    public final static LatencyMetrics globalWriteLatency = new 
LatencyMetrics(globalFactory, globalAliasFactory, "Write");
-    public final static LatencyMetrics globalRangeLatency = new 
LatencyMetrics(globalFactory, globalAliasFactory, "Range");
-
     /** When sampler activated, will track the most frequently read partitions 
**/
     public final Sampler<ByteBuffer> topReadPartitionFrequency;
     /** When sampler activated, will track the most frequently written to 
partitions **/
@@ -284,7 +290,7 @@ public class TableMetrics
         return Pair.create(filtered, total);
     }
 
-    public static final Gauge<Double> globalPercentRepaired = 
Metrics.register(globalFactory.createMetricName("PercentRepaired"),
+    public static final Gauge<Double> globalPercentRepaired = 
Metrics.register(GLOBAL_FACTORY.createMetricName("PercentRepaired"),
                                                                                
new Gauge<Double>()
     {
         public Double getValue()
@@ -296,15 +302,15 @@ public class TableMetrics
         }
     });
 
-    public static final Gauge<Long> globalBytesRepaired = 
Metrics.register(globalFactory.createMetricName("BytesRepaired"),
+    public static final Gauge<Long> globalBytesRepaired = 
Metrics.register(GLOBAL_FACTORY.createMetricName("BytesRepaired"),
                                                                            () 
-> totalNonSystemTablesSize(SSTableReader::isRepaired).left);
 
     public static final Gauge<Long> globalBytesUnrepaired = 
-        Metrics.register(globalFactory.createMetricName("BytesUnrepaired"),
+        Metrics.register(GLOBAL_FACTORY.createMetricName("BytesUnrepaired"),
                          () -> totalNonSystemTablesSize(s -> !s.isRepaired() 
&& !s.isPendingRepair()).left);
 
     public static final Gauge<Long> globalBytesPendingRepair = 
-        Metrics.register(globalFactory.createMetricName("BytesPendingRepair"),
+        Metrics.register(GLOBAL_FACTORY.createMetricName("BytesPendingRepair"),
                          () -> 
totalNonSystemTablesSize(SSTableReader::isPendingRepair).left);
 
     public final Meter readRepairRequests;
@@ -321,15 +327,11 @@ public class TableMetrics
     public final Histogram rfpRowsCachedPerQuery;
 
     public final EnumMap<SamplerType, Sampler<?>> samplers;
-    /**
-     * stores metrics that will be rolled into a single global metric
-     */
-    public final static ConcurrentMap<String, Set<Metric>> allTableMetrics = 
Maps.newConcurrentMap();
 
     /**
      * Stores all metrics created that can be used when unregistering
      */
-    public final static Set<ReleasableMetric> all = Sets.newHashSet();
+    private final Set<ReleasableMetric> all = Sets.newHashSet();
 
     private interface GetHistogram
     {
@@ -566,9 +568,9 @@ public class TableMetrics
             }
         });
 
-        readLatency = createLatencyMetrics("Read", 
cfs.keyspace.metric.readLatency, globalReadLatency);
-        writeLatency = createLatencyMetrics("Write", 
cfs.keyspace.metric.writeLatency, globalWriteLatency);
-        rangeLatency = createLatencyMetrics("Range", 
cfs.keyspace.metric.rangeLatency, globalRangeLatency);
+        readLatency = createLatencyMetrics("Read", 
cfs.keyspace.metric.readLatency, GLOBAL_READ_LATENCY);
+        writeLatency = createLatencyMetrics("Write", 
cfs.keyspace.metric.writeLatency, GLOBAL_WRITE_LATENCY);
+        rangeLatency = createLatencyMetrics("Range", 
cfs.keyspace.metric.rangeLatency, GLOBAL_RANGE_LATENCY);
         pendingFlushes = createTableCounter("PendingFlushes");
         bytesFlushed = createTableCounter("BytesFlushed");
 
@@ -605,7 +607,7 @@ public class TableMetrics
             public Long getValue()
             {
                 long min = Long.MAX_VALUE;
-                for (Metric cfGauge : allTableMetrics.get("MinPartitionSize"))
+                for (Metric cfGauge : 
ALL_TABLE_METRICS.get("MinPartitionSize"))
                 {
                     min = Math.min(min, ((Gauge<? extends Number>) 
cfGauge).getValue().longValue());
                 }
@@ -629,7 +631,7 @@ public class TableMetrics
             public Long getValue()
             {
                 long max = 0;
-                for (Metric cfGauge : allTableMetrics.get("MaxPartitionSize"))
+                for (Metric cfGauge : 
ALL_TABLE_METRICS.get("MaxPartitionSize"))
                 {
                     max = Math.max(max, ((Gauge<? extends Number>) 
cfGauge).getValue().longValue());
                 }
@@ -897,7 +899,7 @@ public class TableMetrics
         unleveledSSTables = createTableGauge("UnleveledSSTables", 
cfs::getUnleveledSSTables, () -> {
             // global gauge
             int cnt = 0;
-            for (Metric cfGauge : allTableMetrics.get("UnleveledSSTables"))
+            for (Metric cfGauge : ALL_TABLE_METRICS.get("UnleveledSSTables"))
             {
                 cnt += ((Gauge<? extends Number>) 
cfGauge).getValue().intValue();
             }
@@ -944,7 +946,7 @@ public class TableMetrics
         Gauge<T> cfGauge = Metrics.register(factory.createMetricName(name), 
aliasFactory.createMetricName(alias), gauge);
         if (register(name, alias, cfGauge) && globalGauge != null)
         {
-            Metrics.register(globalFactory.createMetricName(name), 
globalAliasFactory.createMetricName(alias), globalGauge);
+            Metrics.register(GLOBAL_FACTORY.createMetricName(name), 
GLOBAL_ALIAS_FACTORY.createMetricName(alias), globalGauge);
         }
         return cfGauge;
     }
@@ -969,11 +971,11 @@ public class TableMetrics
         
         if (register(name, name, deprecated, cfGauge))
         {
-            Metrics.register(globalFactory.createMetricName(name), 
+            Metrics.register(GLOBAL_FACTORY.createMetricName(name),
                              globalGauge,
-                             globalAliasFactory.createMetricName(name),
-                             globalFactory.createMetricName(deprecated),
-                             globalAliasFactory.createMetricName(deprecated));
+                             GLOBAL_ALIAS_FACTORY.createMetricName(name),
+                             GLOBAL_FACTORY.createMetricName(deprecated),
+                             
GLOBAL_ALIAS_FACTORY.createMetricName(deprecated));
         }
         return cfGauge;
     }
@@ -992,14 +994,14 @@ public class TableMetrics
         Counter cfCounter = Metrics.counter(factory.createMetricName(name), 
aliasFactory.createMetricName(alias));
         if (register(name, alias, cfCounter))
         {
-            Metrics.register(globalFactory.createMetricName(name),
-                             globalAliasFactory.createMetricName(alias),
+            Metrics.register(GLOBAL_FACTORY.createMetricName(name),
+                             GLOBAL_ALIAS_FACTORY.createMetricName(alias),
                              new Gauge<Long>()
             {
                 public Long getValue()
                 {
                     long total = 0;
-                    for (Metric cfGauge : allTableMetrics.get(name))
+                    for (Metric cfGauge : ALL_TABLE_METRICS.get(name))
                     {
                         total += ((Counter) cfGauge).getCount();
                     }
@@ -1070,8 +1072,8 @@ public class TableMetrics
         register(name, alias, cfHistogram);
         return new TableHistogram(cfHistogram,
                                   keyspaceHistogram,
-                                  
Metrics.histogram(globalFactory.createMetricName(name),
-                                                    
globalAliasFactory.createMetricName(alias),
+                                  
Metrics.histogram(GLOBAL_FACTORY.createMetricName(name),
+                                                    
GLOBAL_ALIAS_FACTORY.createMetricName(alias),
                                                     considerZeroes));
     }
 
@@ -1091,7 +1093,7 @@ public class TableMetrics
     {
         Timer cfTimer = Metrics.timer(factory.createMetricName(name), 
aliasFactory.createMetricName(name));
         register(name, name, keyspaceTimer);
-        Timer global = Metrics.timer(globalFactory.createMetricName(name), 
globalAliasFactory.createMetricName(name));
+        Timer global = Metrics.timer(GLOBAL_FACTORY.createMetricName(name), 
GLOBAL_ALIAS_FACTORY.createMetricName(name));
 
         return new TableTimer(cfTimer, keyspaceTimer, global);
     }
@@ -1114,8 +1116,8 @@ public class TableMetrics
         register(name, alias, meter);
         return new TableMeter(meter,
                               keyspaceMeter,
-                              
Metrics.meter(globalFactory.createMetricName(name),
-                                            
globalAliasFactory.createMetricName(alias)));
+                              
Metrics.meter(GLOBAL_FACTORY.createMetricName(name),
+                                            
GLOBAL_ALIAS_FACTORY.createMetricName(alias)));
     }
 
     private LatencyMetrics createLatencyMetrics(String namePrefix, 
LatencyMetrics ... parents)
@@ -1145,8 +1147,8 @@ public class TableMetrics
      */
     private boolean register(String name, String alias, String deprecated, 
Metric metric)
     {
-        boolean ret = allTableMetrics.putIfAbsent(name, 
ConcurrentHashMap.newKeySet()) == null;
-        allTableMetrics.get(name).add(metric);
+        boolean ret = ALL_TABLE_METRICS.putIfAbsent(name, 
ConcurrentHashMap.newKeySet()) == null;
+        ALL_TABLE_METRICS.get(name).add(metric);
         all.add(() -> releaseMetric(name, alias, deprecated));
         return ret;
     }
@@ -1159,7 +1161,7 @@ public class TableMetrics
         if (metric != null)
         {
             // Metric will be null if we are releasing a view metric.  Views 
have null for ViewLockAcquireTime and ViewLockReadTime
-            allTableMetrics.get(tableMetricName).remove(metric);
+            ALL_TABLE_METRICS.get(tableMetricName).remove(metric);
             CassandraMetricsRegistry.MetricName cfAlias = 
aliasFactory.createMetricName(cfMetricName);
             
             if (tableMetricAlias != null)
@@ -1331,7 +1333,7 @@ public class TableMetrics
         public Long getValue()
         {
             long total = 0;
-            for (Metric cfGauge : allTableMetrics.get(name))
+            for (Metric cfGauge : ALL_TABLE_METRICS.get(name))
             {
                 total = total + ((Gauge<? extends Number>) 
cfGauge).getValue().longValue();
             }
diff --git a/src/java/org/apache/cassandra/utils/MBeanWrapper.java 
b/src/java/org/apache/cassandra/utils/MBeanWrapper.java
index 81d8dae..f8bc439 100644
--- a/src/java/org/apache/cassandra/utils/MBeanWrapper.java
+++ b/src/java/org/apache/cassandra/utils/MBeanWrapper.java
@@ -21,14 +21,14 @@ package org.apache.cassandra.utils;
 import java.lang.management.ManagementFactory;
 import java.util.function.Consumer;
 import javax.management.MBeanServer;
+import javax.management.MalformedObjectNameException;
 import javax.management.ObjectName;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import org.apache.cassandra.config.CassandraRelevantProperties;
-
 import static 
org.apache.cassandra.config.CassandraRelevantProperties.IS_DISABLED_MBEAN_REGISTRATION;
+import static 
org.apache.cassandra.config.CassandraRelevantProperties.MBEAN_REGISTRATION_CLASS;
 
 /**
  * Helper class to avoid catching and rethrowing checked exceptions on MBean 
and
@@ -38,9 +38,18 @@ public interface MBeanWrapper
 {
     static final Logger logger = LoggerFactory.getLogger(MBeanWrapper.class);
 
-    static final MBeanWrapper instance = 
IS_DISABLED_MBEAN_REGISTRATION.getBoolean() ?
-                                         new NoOpMBeanWrapper() :
-                                         new PlatformMBeanWrapper();
+    static final MBeanWrapper instance = create();
+
+    static MBeanWrapper create()
+    {
+        if (IS_DISABLED_MBEAN_REGISTRATION.getBoolean())
+            return new NoOpMBeanWrapper();
+
+        String klass = MBEAN_REGISTRATION_CLASS.getString();
+        if (klass == null)
+            return new PlatformMBeanWrapper();
+        return FBUtilities.construct(klass, "mbean");
+    }
 
     // Passing true for graceful will log exceptions instead of rethrowing them
     public void registerMBean(Object obj, ObjectName mbeanName, OnException 
onException);
@@ -49,7 +58,13 @@ public interface MBeanWrapper
         registerMBean(obj, mbeanName, OnException.THROW);
     }
 
-    public void registerMBean(Object obj, String mbeanName, OnException 
onException);
+    default void registerMBean(Object obj, String mbeanName, OnException 
onException)
+    {
+        ObjectName name = create(mbeanName, onException);
+        if (name == null)
+            return;
+        registerMBean(obj, name, onException);
+    }
     default void registerMBean(Object obj, String mbeanName)
     {
         registerMBean(obj, mbeanName, OnException.THROW);
@@ -61,7 +76,13 @@ public interface MBeanWrapper
         return isRegistered(mbeanName, OnException.THROW);
     }
 
-    public boolean isRegistered(String mbeanName, OnException onException);
+    default boolean isRegistered(String mbeanName, OnException onException)
+    {
+        ObjectName name = create(mbeanName, onException);
+        if (name == null)
+            return false;
+        return isRegistered(name, onException);
+    }
     default boolean isRegistered(String mbeanName)
     {
         return isRegistered(mbeanName, OnException.THROW);
@@ -73,12 +94,31 @@ public interface MBeanWrapper
         unregisterMBean(mbeanName, OnException.THROW);
     }
 
-    public void unregisterMBean(String mbeanName, OnException onException);
+    default void unregisterMBean(String mbeanName, OnException onException)
+    {
+        ObjectName name = create(mbeanName, onException);
+        if (name == null)
+            return;
+        unregisterMBean(name, onException);
+    }
     default void unregisterMBean(String mbeanName)
     {
         unregisterMBean(mbeanName, OnException.THROW);
     }
 
+    static ObjectName create(String mbeanName, OnException onException)
+    {
+        try
+        {
+            return new ObjectName(mbeanName);
+        }
+        catch (MalformedObjectNameException e)
+        {
+            onException.handler.accept(e);
+            return null;
+        }
+    }
+
     static class NoOpMBeanWrapper implements MBeanWrapper
     {
         public void registerMBean(Object obj, ObjectName mbeanName, 
OnException onException) {}
@@ -104,18 +144,6 @@ public interface MBeanWrapper
             }
         }
 
-        public void registerMBean(Object obj, String mbeanName, OnException 
onException)
-        {
-            try
-            {
-                mbs.registerMBean(obj, new ObjectName(mbeanName));
-            }
-            catch (Exception e)
-            {
-                onException.handler.accept(e);
-            }
-        }
-
         public boolean isRegistered(ObjectName mbeanName, OnException 
onException)
         {
             try
@@ -129,19 +157,6 @@ public interface MBeanWrapper
             return false;
         }
 
-        public boolean isRegistered(String mbeanName, OnException onException)
-        {
-            try
-            {
-                return mbs.isRegistered(new ObjectName(mbeanName));
-            }
-            catch (Exception e)
-            {
-                onException.handler.accept(e);
-            }
-            return false;
-        }
-
         public void unregisterMBean(ObjectName mbeanName, OnException 
onException)
         {
             try
@@ -153,18 +168,6 @@ public interface MBeanWrapper
                 onException.handler.accept(e);
             }
         }
-
-        public void unregisterMBean(String mbeanName, OnException onException)
-        {
-            try
-            {
-                mbs.unregisterMBean(new ObjectName(mbeanName));
-            }
-            catch (Exception e)
-            {
-                onException.handler.accept(e);
-            }
-        }
     }
 
     public enum OnException
@@ -173,7 +176,7 @@ public interface MBeanWrapper
         LOG(e -> { logger.error("Error in MBean wrapper: ", e); }),
         IGNORE(e -> {});
 
-        private Consumer<Exception> handler;
+        public final Consumer<Exception> handler;
         OnException(Consumer<Exception> handler)
         {
             this.handler = handler;
diff --git 
a/test/distributed/org/apache/cassandra/distributed/impl/AbstractCluster.java 
b/test/distributed/org/apache/cassandra/distributed/impl/AbstractCluster.java
index 228385e..1ae594f 100644
--- 
a/test/distributed/org/apache/cassandra/distributed/impl/AbstractCluster.java
+++ 
b/test/distributed/org/apache/cassandra/distributed/impl/AbstractCluster.java
@@ -305,6 +305,7 @@ public abstract class AbstractCluster<I extends IInstance> 
implements ICluster<I
             }
         }
 
+        @Override
         public void uncaughtException(Thread thread, Throwable throwable)
         {
             IInvokableInstance delegate = this.delegate;
@@ -313,6 +314,13 @@ public abstract class AbstractCluster<I extends IInstance> 
implements ICluster<I
             else
                 logger.error("uncaught exception in thread {}", thread, 
throwable);
         }
+
+        @Override
+        public String toString()
+        {
+            IInvokableInstance delegate = this.delegate;
+            return delegate == null ? "node" + config.num() : 
delegate.toString();
+        }
     }
 
     protected AbstractCluster(AbstractBuilder<I, ? extends ICluster<I>, ?> 
builder)
diff --git 
a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java 
b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
index 2a71ec2..8dd13d7 100644
--- a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
+++ b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
@@ -489,11 +489,10 @@ public class Instance extends IsolatedExecutor implements 
IInvokableInstance
                             GossipHelper.unsafeStatusToNormal(this, 
(IInstance) peer);
                     });
 
+                    StorageService.instance.setUpDistributedSystemKeyspaces();
                     StorageService.instance.setNormalModeUnsafe();
                 }
 
-                StorageService.instance.ensureTraceKeyspace();
-
                 // Populate tokenMetadata for the second time,
                 // see org.apache.cassandra.service.CassandraDaemon.setup
                 StorageService.instance.populateTokenMetadata();
@@ -737,6 +736,12 @@ public class Instance extends IsolatedExecutor implements 
IInvokableInstance
         }).call();
     }
 
+    @Override
+    public String toString()
+    {
+        return "node" + config.num();
+    }
+
     private static class CapturingOutput implements Closeable
     {
         @SuppressWarnings("resource")
diff --git 
a/test/distributed/org/apache/cassandra/distributed/test/metric/TableMetricTest.java
 
b/test/distributed/org/apache/cassandra/distributed/test/metric/TableMetricTest.java
new file mode 100644
index 0000000..751a8e3
--- /dev/null
+++ 
b/test/distributed/org/apache/cassandra/distributed/test/metric/TableMetricTest.java
@@ -0,0 +1,272 @@
+/*
+ * 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.cassandra.distributed.test.metric;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.stream.Collectors;
+import javax.management.InstanceAlreadyExistsException;
+import javax.management.InstanceNotFoundException;
+import javax.management.ObjectName;
+
+import com.google.common.collect.ImmutableSet;
+import org.junit.Assert;
+import org.junit.Test;
+
+import org.apache.cassandra.auth.AuthKeyspace;
+import org.apache.cassandra.db.ColumnFamilyStore;
+import org.apache.cassandra.db.SystemKeyspace;
+import org.apache.cassandra.distributed.Cluster;
+import org.apache.cassandra.distributed.api.IInvokableInstance;
+import org.apache.cassandra.distributed.test.TestBaseImpl;
+import org.apache.cassandra.repair.SystemDistributedKeyspace;
+import org.apache.cassandra.schema.Schema;
+import org.apache.cassandra.schema.SchemaKeyspace;
+import org.apache.cassandra.tracing.TraceKeyspace;
+import org.apache.cassandra.utils.MBeanWrapper;
+
+import static 
org.apache.cassandra.config.CassandraRelevantProperties.IS_DISABLED_MBEAN_REGISTRATION;
+import static 
org.apache.cassandra.config.CassandraRelevantProperties.MBEAN_REGISTRATION_CLASS;
+
+public class TableMetricTest extends TestBaseImpl
+{
+    static
+    {
+        MBEAN_REGISTRATION_CLASS.setString(MapMBeanWrapper.class.getName());
+        IS_DISABLED_MBEAN_REGISTRATION.setBoolean(false);
+    }
+
+    private static volatile Map<String, Collection<String>> SYSTEM_TABLES = 
null;
+    private static Set<String> TABLE_METRIC_NAMES = 
ImmutableSet.of("WriteLatency");
+
+    /**
+     * Makes sure that all system tables have the expected metrics
+     * @throws IOException
+     */
+    @Test
+    public void systemTables() throws IOException
+    {
+        try (Cluster cluster = Cluster.build(1).start())
+        {
+            loadSystemTables(cluster);
+            assertSystemTableMetrics(cluster);
+        }
+    }
+
+    /**
+     * Tests that other table metrics are not modified when a single table is 
modified/deleted.
+     *
+     * @see <a 
href="https://issues.apache.org/jira/browse/CASSANDRA-16095";>CASSANDRA-16095</a>
+     */
+    @Test
+    public void userTables() throws IOException
+    {
+        try (Cluster cluster = init(Cluster.build(3).start()))
+        {
+            loadSystemTables(cluster);
+            assertSystemTableMetrics(cluster);
+
+            cluster.schemaChange(withKeyspace("CREATE TABLE %s.tbl (pk bigint 
PRIMARY KEY)"));
+            cluster.forEach(i -> assertTableMetricsExist(i, KEYSPACE, "tbl"));
+
+            // alter table can change metrics, so monitor for it
+            cluster.schemaChange(withKeyspace("ALTER TABLE %s.tbl WITH comment 
= 'testing'"));
+            cluster.forEach(i -> assertTableMetricsExist(i, KEYSPACE, "tbl"));
+
+            cluster.schemaChange(withKeyspace("ALTER TABLE %s.tbl ADD (value 
bigint)"));
+            cluster.forEach(i -> assertTableMetricsExist(i, KEYSPACE, "tbl"));
+
+            cluster.schemaChange(withKeyspace("ALTER TABLE %s.tbl RENAME pk TO 
pk2"));
+            cluster.forEach(i -> assertTableMetricsExist(i, KEYSPACE, "tbl"));
+
+            cluster.schemaChange(withKeyspace("ALTER TABLE %s.tbl DROP 
value"));
+            cluster.forEach(i -> assertTableMetricsExist(i, KEYSPACE, "tbl"));
+
+            // drop and make sure table no longer exists
+            cluster.schemaChange(withKeyspace("DROP TABLE %s.tbl"));
+            cluster.forEach(i -> assertTableMetricsDoesNotExist(i, KEYSPACE, 
"tbl"));
+
+            cluster.schemaChange(withKeyspace("DROP KEYSPACE %s"));
+            cluster.forEach(i -> assertKeyspaceMetricDoesNotExists(i, 
KEYSPACE));
+
+            // no other table impacted?
+            assertSystemTableMetrics(cluster);
+        }
+    }
+
+    private static void loadSystemTables(Cluster cluster)
+    {
+        SYSTEM_TABLES = cluster.get(1).callOnInstance(() -> {
+            Map<String, Collection<String>> map = new HashMap<>();
+            Arrays.asList(SystemKeyspace.metadata(), AuthKeyspace.metadata(), 
SystemDistributedKeyspace.metadata(),
+                          SchemaKeyspace.metadata(), TraceKeyspace.metadata())
+                  .forEach(meta -> {
+                      Set<String> tables = meta.tables.stream().map(t -> 
t.name).collect(Collectors.toSet());
+                      map.put(meta.name, tables);
+                  });
+            return map;
+        });
+    }
+
+    private static void assertSystemTableMetrics(Cluster cluster)
+    {
+        for (String keyspace : SYSTEM_TABLES.keySet())
+        {
+            for (String table : SYSTEM_TABLES.get(keyspace))
+            {
+                cluster.forEach(i -> assertTableMetricsExist(i, keyspace, 
table));
+            }
+        }
+    }
+
+    private static void assertTableMetricsExist(IInvokableInstance inst, 
String keyspace, String table)
+    {
+        assertTableMBeanExists(inst, keyspace, table);
+        for (String metric : TABLE_METRIC_NAMES)
+            assertTableMetricExists(inst, keyspace, table, metric);
+    }
+
+    private static void assertTableMetricsDoesNotExist(IInvokableInstance 
inst, String keyspace, String table)
+    {
+        assertTableMBeanDoesNotExists(inst, keyspace, table);
+        for (String metric : TABLE_METRIC_NAMES)
+            assertTableMetricDoesNotExists(inst, keyspace, table, metric);
+    }
+
+    private static void assertKeyspaceMetricDoesNotExists(IInvokableInstance 
inst, String keyspace)
+    {
+        for (String metric : TABLE_METRIC_NAMES)
+            assertKeyspaceMetricDoesNotExists(inst, keyspace, metric);
+    }
+
+    private static void assertTableMBeanExists(IInvokableInstance inst, String 
keyspace, String table)
+    {
+        inst.runOnInstance(() -> {
+            // cast only to make sure it linked properly
+            MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
+            Assert.assertTrue("Unable to find table mbean for " + keyspace + 
"." + table,
+                              
mbeans.isRegistered(ColumnFamilyStore.getTableMBeanName(keyspace, table, 
false)));
+            Assert.assertTrue("Unable to find column family mbean for " + 
keyspace + "." + table,
+                              
mbeans.isRegistered(ColumnFamilyStore.getColumnFamilieMBeanName(keyspace, 
table, false)));
+        });
+    }
+
+    private static void assertTableMBeanDoesNotExists(IInvokableInstance inst, 
String keyspace, String table)
+    {
+        inst.runOnInstance(() -> {
+            // cast only to make sure it linked properly
+            MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
+            Assert.assertFalse("Found table mbean for " + keyspace + "." + 
table,
+                               
mbeans.isRegistered(ColumnFamilyStore.getTableMBeanName(keyspace, table, 
false)));
+            Assert.assertFalse("Found column family mbean for " + keyspace + 
"." + table,
+                               
mbeans.isRegistered(ColumnFamilyStore.getColumnFamilieMBeanName(keyspace, 
table, false)));
+        });
+    }
+
+    private static void assertTableMetricExists(IInvokableInstance inst, 
String keyspace, String table, String name)
+    {
+        inst.runOnInstance(() -> {
+            // cast only to make sure it linked properly
+            MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
+            String mbean = getTableMetricName(keyspace, table, name);
+            Assert.assertTrue("Unable to find metric " + name + " for " + 
keyspace + "." + table, mbeans.isRegistered(mbean));
+
+            // verify replicated to keyspace
+            String keyspaceMBean = getKeyspaceMetricName(keyspace, name);
+            Assert.assertTrue("Unable to find keyspace metric " + 
keyspaceMBean + " for " + keyspace, mbeans.isRegistered(keyspaceMBean));
+        });
+    }
+
+    private static void assertTableMetricDoesNotExists(IInvokableInstance 
inst, String keyspace, String table, String name)
+    {
+        inst.runOnInstance(() -> {
+            // cast only to make sure it linked properly
+            MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
+            String mbean = getTableMetricName(keyspace, table, name);
+            Assert.assertFalse("Found metric " + name + " for " + keyspace + 
"." + table, mbeans.isRegistered(mbean));
+
+            // validate keyspace metric
+            assertKeyspaceMetricMayExists(mbeans, keyspace, name);
+        });
+    }
+
+    private static void assertKeyspaceMetricMayExists(MapMBeanWrapper mbeans, 
String keyspace, String name)
+    {
+        String keyspaceMBean = getKeyspaceMetricName(keyspace, name);
+        boolean keyspaceExists = Schema.instance.getKeyspaceMetadata(keyspace) 
!= null;
+        String errorMessage = keyspaceExists ?
+                              "Unable to find keyspace metric " + 
keyspaceMBean + " for " + keyspace :
+                              "Found keyspace metric " + keyspaceMBean + " for 
" + keyspace;
+        Assert.assertEquals(errorMessage, keyspaceExists, 
mbeans.isRegistered(keyspaceMBean));
+    }
+
+    private static void assertKeyspaceMetricDoesNotExists(IInvokableInstance 
inst, String keyspace, String name)
+    {
+        inst.runOnInstance(() -> {
+            // cast only to make sure it linked properly
+            MapMBeanWrapper mbeans = (MapMBeanWrapper) MBeanWrapper.instance;
+
+            String keyspaceMBean = getKeyspaceMetricName(keyspace, name);
+            Assert.assertFalse("Found keyspace metric " + keyspaceMBean + " 
for " + keyspace, mbeans.isRegistered(keyspaceMBean));
+        });
+    }
+
+    private static String getKeyspaceMetricName(String keyspace, String name)
+    {
+        return 
String.format("org.apache.cassandra.metrics:type=Keyspace,keyspace=%s,name=%s", 
keyspace, name);
+    }
+
+    private static String getTableMetricName(String keyspace, String table, 
String name)
+    {
+        return 
String.format("org.apache.cassandra.metrics:type=Table,keyspace=%s,scope=%s,name=%s",
 keyspace, table, name);
+    }
+
+    public static final class MapMBeanWrapper implements MBeanWrapper
+    {
+        private final ConcurrentMap<ObjectName, Object> map = new 
ConcurrentHashMap<>();
+
+        @Override
+        public void registerMBean(Object obj, ObjectName mbeanName, 
OnException onException)
+        {
+            Object current = map.putIfAbsent(mbeanName, obj);
+            if (current != null)
+                onException.handler.accept(new 
InstanceAlreadyExistsException("MBean " + mbeanName + " already exists"));
+        }
+
+        @Override
+        public boolean isRegistered(ObjectName mbeanName, OnException 
onException)
+        {
+            return map.containsKey(mbeanName);
+        }
+
+        @Override
+        public void unregisterMBean(ObjectName mbeanName, OnException 
onException)
+        {
+            Object previous = map.remove(mbeanName);
+            if (previous == null)
+                onException.handler.accept(new 
InstanceNotFoundException("MBean " + mbeanName + " was not found"));
+        }
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to