Add ResourceWatcher to generalize reloading Strategy configurations. patch by Jon Hermes and jbellis for CASSANDRA-1105

git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@948829 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
Jonathan Ellis 2010-05-27 13:28:42 +00:00
parent 8c0c3ebaae
commit c924ce55eb
5 changed files with 114 additions and 123 deletions

View File

@ -1,4 +1,5 @@
package org.apache.cassandra.locator;
/*
*
* Licensed to the Apache Software Foundation (ASF) under one
@ -21,11 +22,7 @@ package org.apache.cassandra.locator;
*/
import java.io.FileReader;
import java.io.IOException;
import java.io.IOError;
import java.net.InetAddress;
import java.net.URL;
import java.util.*;
import java.util.Map.Entry;
@ -33,9 +30,11 @@ import com.google.common.collect.Multimap;
import org.apache.cassandra.config.ConfigurationException;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.utils.ResourceWatcher;
import org.apache.cassandra.service.*;
import org.apache.cassandra.thrift.ConsistencyLevel;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.WrappedRunnable;
/**
* This Replication Strategy takes a property file that gives the intended
@ -63,38 +62,36 @@ public class DatacenterShardStrategy extends AbstractReplicationStrategy
if ((!(snitch instanceof AbstractRackAwareSnitch)))
throw new IllegalArgumentException("DatacenterShardStrategy requires a rack-aware endpointsnitch");
this.snitch = (AbstractRackAwareSnitch)snitch;
ClassLoader loader = PropertyFileSnitch.class.getClassLoader();
URL scpurl = loader.getResource(DATACENTER_PROPERTY_FILENAME);
if (scpurl == null)
reloadConfiguration();
Runnable runnable = new WrappedRunnable()
{
throw new RuntimeException("unable to locate " + DATACENTER_PROPERTY_FILENAME);
}
String dcPropertyFile = scpurl.getFile();
try
{
Properties props = new Properties();
props.load(new FileReader(dcPropertyFile));
for (Object key : props.keySet())
protected void runMayThrow() throws ConfigurationException
{
String[] keys = ((String)key).split(":");
Map<String, Integer> map = datacenters.get(keys[0]);
if (null == map)
{
map = new HashMap<String, Integer>();
}
map.put(keys[1], Integer.parseInt((String)props.get(key)));
datacenters.put(keys[0], map);
reloadConfiguration();
}
}
catch (IOException ioe)
{
throw new IOError(ioe);
}
};
ResourceWatcher.watch(DATACENTER_PROPERTY_FILENAME, runnable, 60 * 1000);
loadEndpoints(tokenMetadata);
}
public void reloadConfiguration() throws ConfigurationException
{
Properties props = PropertyFileSnitch.resourceToProperties(DATACENTER_PROPERTY_FILENAME);
for (Object key : props.keySet())
{
String[] keys = ((String)key).split(":");
Map<String, Integer> map = datacenters.get(keys[0]);
if (null == map)
{
map = new HashMap<String, Integer>();
}
map.put(keys[1], Integer.parseInt((String)props.get(key)));
datacenters.put(keys[0], map);
}
}
private synchronized void loadEndpoints(TokenMetadata metadata) throws ConfigurationException
{
String localDC = snitch.getDatacenter(DatabaseDescriptor.getListenAddress());

View File

@ -20,18 +20,14 @@ package org.apache.cassandra.locator;
import java.io.FileReader;
import java.io.IOException;
import java.lang.management.ManagementFactory;
import java.net.UnknownHostException;
import java.net.URL;
import java.net.InetAddress;
import java.util.Properties;
import java.util.StringTokenizer;
import javax.management.MBeanServer;
import javax.management.ObjectName;
import java.net.InetAddress;
import org.apache.cassandra.config.ConfigurationException;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.ResourceWatcher;
import org.apache.cassandra.utils.WrappedRunnable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -40,11 +36,11 @@ import org.slf4j.LoggerFactory;
* <p/>
* Based on a properties file configuration.
*/
public class PropertyFileSnitch extends AbstractRackAwareSnitch implements PropertyFileSnitchMBean {
public class PropertyFileSnitch extends AbstractRackAwareSnitch {
/**
* A list of properties with keys being host:port and values being datacenter:rack
*/
private Properties hostProperties = new Properties();
private volatile Properties hostProperties = new Properties();
/**
* The default rack property file to be read.
@ -56,18 +52,16 @@ public class PropertyFileSnitch extends AbstractRackAwareSnitch implements Prope
*/
private static Logger logger_ = LoggerFactory.getLogger(PropertyFileSnitch.class);
public PropertyFileSnitch() throws ConfigurationException
{
public PropertyFileSnitch() throws ConfigurationException {
reloadConfiguration();
try
Runnable runnable = new WrappedRunnable()
{
MBeanServer mbs = ManagementFactory.getPlatformMBeanServer();
mbs.registerMBean(this, new ObjectName(MBEAN_OBJECT_NAME));
}
catch (Exception e)
{
throw new RuntimeException(e);
}
protected void runMayThrow() throws ConfigurationException
{
reloadConfiguration();
}
};
ResourceWatcher.watch(RACK_PROPERTY_FILENAME, runnable, 60 * 1000);
}
/**
@ -115,40 +109,25 @@ public class PropertyFileSnitch extends AbstractRackAwareSnitch implements Prope
return getEndpointInfo(endpoint)[1];
}
/**
* @return the <tt>String</tt> representation of the configuration
*/
public String displayConfiguration() {
StringBuffer configurationString = new StringBuffer("Current rack configuration\n=================\n");
for (Object key: hostProperties.keySet()) {
String endpoint = (String) key;
String value = hostProperties.getProperty(endpoint);
configurationString.append(endpoint).append("=").append(value).append("\n");
}
return configurationString.toString();
}
/**
* Reloads the configuration from the file
*/
public void reloadConfiguration() throws ConfigurationException
{
ClassLoader loader = PropertyFileSnitch.class.getClassLoader();
URL scpurl = loader.getResource(RACK_PROPERTY_FILENAME);
if (scpurl == null)
throw new ConfigurationException("unable to locate " + RACK_PROPERTY_FILENAME);
hostProperties = resourceToProperties(RACK_PROPERTY_FILENAME);
}
String rackPropertyFilename = scpurl.getFile();
public static Properties resourceToProperties(String filename) throws ConfigurationException
{
String rackPropertyFilename = FBUtilities.resourceToFile(filename);
Properties localHostProperties;
try
{
Properties localHostProperties = new Properties();
localHostProperties = new Properties();
localHostProperties.load(new FileReader(rackPropertyFilename));
hostProperties = localHostProperties;
}
catch (IOException ioe)
catch (IOException e)
{
throw new ConfigurationException("Could not process " + rackPropertyFilename, ioe);
throw new ConfigurationException("Unable to load " + rackPropertyFilename, e);
}
return localHostProperties;
}
}

View File

@ -1,49 +0,0 @@
/**
* 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.locator;
import java.io.IOException;
import org.apache.cassandra.config.ConfigurationException;
/**
* PropertyFileSnitchMBean
*
* PropertyFileSnitchMBean is the management interface for Digg's EndpointSnitch MBean
*
* @author Sammy Yu <syu@sammyyu.net>
*
*/
public interface PropertyFileSnitchMBean {
/**
* The object name of the mbean.
*/
public static String MBEAN_OBJECT_NAME = "org.apache.cassandra.locator:type=EndpointSnitch";
/**
* Reload the rack configuration
*/
public void reloadConfiguration() throws ConfigurationException;
/**
* Display the current configuration
*/
public String displayConfiguration();
}

View File

@ -21,6 +21,7 @@ package org.apache.cassandra.utils;
import java.io.*;
import java.math.BigInteger;
import java.net.InetAddress;
import java.net.URL;
import java.net.UnknownHostException;
import java.nio.ByteBuffer;
import java.nio.charset.Charset;
@ -33,6 +34,8 @@ import java.util.zip.DataFormatException;
import java.util.zip.Deflater;
import java.util.zip.Inflater;
import org.apache.cassandra.config.ConfigurationException;
import org.apache.cassandra.locator.PropertyFileSnitch;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -548,4 +551,14 @@ public class FBUtilities
ByteBuffer.wrap(bytes).putLong(n);
return bytes;
}
public static String resourceToFile(String filename) throws ConfigurationException
{
ClassLoader loader = PropertyFileSnitch.class.getClassLoader();
URL scpurl = loader.getResource(filename);
if (scpurl == null)
throw new ConfigurationException("unable to locate " + filename);
return scpurl.getFile();
}
}

View File

@ -0,0 +1,51 @@
package org.apache.cassandra.utils;
import java.io.File;
import java.util.Timer;
import java.util.TimerTask;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class ResourceWatcher
{
private static Timer timer = new Timer("RESOURCE-WATCHER");
public static void watch(String resource, Runnable callback, int period)
{
timer.schedule(new WatchedResource(resource, callback), period, period);
}
public static class WatchedResource extends TimerTask
{
private static Logger logger = LoggerFactory.getLogger(WatchedResource.class);
private String resource;
private Runnable callback;
private long lastLoaded;
public WatchedResource(String resource, Runnable callback)
{
this.resource = resource;
this.callback = callback;
lastLoaded = 0;
}
public void run()
{
try
{
String filename = FBUtilities.resourceToFile(resource);
long lastModified = new File(filename).lastModified();
if (lastModified > lastLoaded)
{
callback.run();
lastLoaded = lastModified;
}
}
catch (Throwable t)
{
logger.error(String.format("Timed run of %s failed.", callback.getClass()), t);
}
}
}
}