diff --git a/CHANGES.txt b/CHANGES.txt index 496885bff0..37839fb461 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -5,6 +5,7 @@ dev * avoid polluting page cache with commitlog or sstable writes and seq scan operations (CASSANDRA-1470) * add RMI authentication options to nodetool (CASSANDRA-1921) + * Make snitches configurable at runtime (CASSANDRA-1374) 0.7.0-rc4 diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 04957e220f..f1f23f07ad 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -400,7 +400,7 @@ public class DatabaseDescriptor IEndpointSnitch snitch = FBUtilities.construct(endpointSnitchClassName, "snitch"); return conf.dynamic_snitch ? new DynamicEndpointSnitch(snitch) : snitch; } - + /** load keyspace (table) definitions, but do not initialize the table instances. */ public static void loadSchemas() throws IOException { @@ -705,6 +705,10 @@ public class DatabaseDescriptor { return snitch; } + public static void setEndpointSnitch(IEndpointSnitch eps) + { + snitch = eps; + } public static IRequestScheduler getRequestScheduler() { @@ -1104,14 +1108,26 @@ public class DatabaseDescriptor { return conf.dynamic_snitch_update_interval_in_ms; } + public static void setDynamicUpdateInterval(Integer dynamicUpdateInterval) + { + conf.dynamic_snitch_update_interval_in_ms = dynamicUpdateInterval; + } public static int getDynamicResetInterval() { return conf.dynamic_snitch_reset_interval_in_ms; } + public static void setDynamicResetInterval(Integer dynamicResetInterval) + { + conf.dynamic_snitch_reset_interval_in_ms = dynamicResetInterval; + } public static double getDynamicBadnessThreshold() { return conf.dynamic_snitch_badness_threshold; } + public static void setDynamicBadnessThreshold(Double dynamicBadnessThreshold) + { + conf.dynamic_snitch_badness_threshold = dynamicBadnessThreshold; + } } diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java index c4e58dc6f6..80a07cbc38 100644 --- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java +++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java @@ -46,9 +46,10 @@ public abstract class AbstractReplicationStrategy private static final Logger logger = LoggerFactory.getLogger(AbstractReplicationStrategy.class); public final String table; - private final TokenMetadata tokenMetadata; - public final IEndpointSnitch snitch; public final Map configOptions; + private final TokenMetadata tokenMetadata; + + public IEndpointSnitch snitch; AbstractReplicationStrategy(String table, TokenMetadata tokenMetadata, IEndpointSnitch snitch, Map configOptions) { diff --git a/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java b/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java index 176e055c2b..f288fdf4a3 100644 --- a/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java +++ b/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java @@ -40,19 +40,23 @@ import org.apache.cassandra.utils.FBUtilities; public class DynamicEndpointSnitch extends AbstractEndpointSnitch implements ILatencySubscriber, DynamicEndpointSnitchMBean { private static final int UPDATES_PER_INTERVAL = 10000; - private static final int UPDATE_INTERVAL_IN_MS = DatabaseDescriptor.getDynamicUpdateInterval(); - private static final int RESET_INTERVAL_IN_MS = DatabaseDescriptor.getDynamicResetInterval(); - private static final double BADNESS_THRESHOLD = DatabaseDescriptor.getDynamicBadnessThreshold(); private static final int WINDOW_SIZE = 100; + + private int UPDATE_INTERVAL_IN_MS = DatabaseDescriptor.getDynamicUpdateInterval(); + private int RESET_INTERVAL_IN_MS = DatabaseDescriptor.getDynamicResetInterval(); + private double BADNESS_THRESHOLD = DatabaseDescriptor.getDynamicBadnessThreshold(); + private String mbeanName; private boolean registered = false; private final ConcurrentHashMap scores = new ConcurrentHashMap(); private final ConcurrentHashMap windows = new ConcurrentHashMap(); private final AtomicInteger intervalupdates = new AtomicInteger(0); + public final IEndpointSnitch subsnitch; public DynamicEndpointSnitch(IEndpointSnitch snitch) { + mbeanName = "org.apache.cassandra.db:type=DynamicEndpointSnitch,instance="+hashCode(); subsnitch = snitch; Runnable update = new Runnable() { @@ -72,11 +76,28 @@ public class DynamicEndpointSnitch extends AbstractEndpointSnitch implements ILa }; StorageService.scheduledTasks.scheduleWithFixedDelay(update, UPDATE_INTERVAL_IN_MS, UPDATE_INTERVAL_IN_MS, TimeUnit.MILLISECONDS); StorageService.scheduledTasks.scheduleWithFixedDelay(reset, RESET_INTERVAL_IN_MS, RESET_INTERVAL_IN_MS, TimeUnit.MILLISECONDS); + registerMBean(); + } + private void registerMBean() + { MBeanServer mbs = ManagementFactory.getPlatformMBeanServer(); try { - mbs.registerMBean(this, new ObjectName("org.apache.cassandra.db:type=DynamicEndpointSnitch,instance="+hashCode())); + mbs.registerMBean(this, new ObjectName(mbeanName)); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + } + + public void unregisterMBean() + { + MBeanServer mbs = ManagementFactory.getPlatformMBeanServer(); + try + { + mbs.unregisterMBean(new ObjectName(mbeanName)); } catch (Exception e) { @@ -208,6 +229,25 @@ public class DynamicEndpointSnitch extends AbstractEndpointSnitch implements ILa { return scores; } + + public int getUpdateInterval() + { + return UPDATE_INTERVAL_IN_MS; + } + public int getResetInterval() + { + return RESET_INTERVAL_IN_MS; + } + public double getBadnessThreshold() + { + return BADNESS_THRESHOLD; + } + public String getSubsnitchClassName() + { + return subsnitch.getClass().getName(); + } + + } /** a threadsafe version of BoundedStatsDeque+ArrivalWindow with modification for arbitrary times **/ diff --git a/src/java/org/apache/cassandra/locator/DynamicEndpointSnitchMBean.java b/src/java/org/apache/cassandra/locator/DynamicEndpointSnitchMBean.java index 26c579977d..5f69709093 100644 --- a/src/java/org/apache/cassandra/locator/DynamicEndpointSnitchMBean.java +++ b/src/java/org/apache/cassandra/locator/DynamicEndpointSnitchMBean.java @@ -24,4 +24,8 @@ import java.util.Map; public interface DynamicEndpointSnitchMBean { public Map getScores(); + public int getUpdateInterval(); + public int getResetInterval(); + public double getBadnessThreshold(); + public String getSubsnitchClassName(); } diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 9615c477c5..b4eb27c2bc 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -27,11 +27,13 @@ import java.nio.ByteBuffer; import java.util.*; import java.util.concurrent.*; import javax.management.MBeanServer; +import javax.management.MalformedObjectNameException; import javax.management.ObjectName; import com.google.common.base.Charsets; import com.google.common.collect.HashMultimap; import com.google.common.collect.Multimap; +import org.apache.cassandra.locator.*; import org.apache.log4j.Level; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; @@ -49,9 +51,6 @@ import org.apache.cassandra.dht.Token; import org.apache.cassandra.gms.*; import org.apache.cassandra.io.DeletionService; import org.apache.cassandra.io.util.FileUtils; -import org.apache.cassandra.locator.AbstractReplicationStrategy; -import org.apache.cassandra.locator.IEndpointSnitch; -import org.apache.cassandra.locator.TokenMetadata; import org.apache.cassandra.net.IAsyncResult; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -2013,4 +2012,28 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe return Collections.unmodifiableList(tableslist); } + public void updateSnitch(String epSnitchClassName, Boolean dynamic, Integer dynamicUpdateInterval, Integer dynamicResetInterval, Double dynamicBadnessThreshold) throws ConfigurationException + { + IEndpointSnitch oldSnitch = DatabaseDescriptor.getEndpointSnitch(); + + // new snitch registers mbean during construction + IEndpointSnitch newSnitch = FBUtilities.construct(epSnitchClassName, "snitch"); + if (dynamic) + { + DatabaseDescriptor.setDynamicUpdateInterval(dynamicUpdateInterval); + DatabaseDescriptor.setDynamicResetInterval(dynamicResetInterval); + DatabaseDescriptor.setDynamicBadnessThreshold(dynamicBadnessThreshold); + newSnitch = new DynamicEndpointSnitch(newSnitch); + } + + // point snitch references to the new instance + DatabaseDescriptor.setEndpointSnitch(newSnitch); + for (String ks : DatabaseDescriptor.getTables()) + { + Table.open(ks).getReplicationStrategy().snitch = newSnitch; + } + + if (oldSnitch instanceof DynamicEndpointSnitch) + ((DynamicEndpointSnitch)oldSnitch).unregisterMBean(); + } } diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java index 92e402d409..acce2b72c4 100644 --- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java +++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java @@ -259,4 +259,15 @@ public interface StorageServiceMBean public Map getOwnership(); public List getKeyspaces(); + + /** + * Change endpointsnitch class and dynamic-ness (and dynamic attributes) at runtime + * @param epSnitchClassName the canonical path name for a class implementing IEndpointSnitch + * @param dynamic boolean that decides whether dynamicsnitch is used or not + * @param dynamicUpdateInterval integer, in ms (default 100) + * @param dynamicResetInterval integer, in ms (default 600,000) + * @param dynamicBadnessThreshold double, (default 0.0) + * @throws ConfigurationException classname not found on classpath + */ + public void updateSnitch(String epSnitchClassName, Boolean dynamic, Integer dynamicUpdateInterval, Integer dynamicResetInterval, Double dynamicBadnessThreshold) throws ConfigurationException; }