mirror of https://github.com/apache/cassandra
merge from 0.8
git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@1152107 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
commit
d505bdef60
13
CHANGES.txt
13
CHANGES.txt
|
|
@ -26,6 +26,19 @@
|
|||
(CASSANDRA-1951)
|
||||
|
||||
|
||||
0.8.3
|
||||
* add ability to drop local reads/writes that are going to timeout
|
||||
(CASSANDRA-2943)
|
||||
* revamp token removal process, keep gossip states for 3 days (CASSANDRA-2946)
|
||||
* don't accept extra args for 0-arg nodetool commands (CASSANDRA-2740)
|
||||
* log unavailableexception details at debug level (CASSANDRA-2856)
|
||||
* expose data_dir though jmx (CASSANDRA-2770)
|
||||
* don't include tmp files as sstable when create cfs (CASSANDRA-2929)
|
||||
* log Java classpath on startup (CASSANDRA-2895)
|
||||
* keep gossipped version in sync with actual on migration coordinator
|
||||
(CASSANDRA-2946)
|
||||
|
||||
|
||||
0.8.2
|
||||
* CQL:
|
||||
- include only one row per unique key for IN queries (CASSANDRA-2717)
|
||||
|
|
|
|||
10
NEWS.txt
10
NEWS.txt
|
|
@ -7,6 +7,16 @@ Upgrading
|
|||
sstableloader tool instead.
|
||||
|
||||
|
||||
0.8.3
|
||||
=====
|
||||
|
||||
Upgrading
|
||||
---------
|
||||
- Token removal has been revamped. Removing tokens in a mixed cluster with
|
||||
0.8.3 will not work, so the entire cluster will need to be running 0.8.3
|
||||
first, except for the dead node.
|
||||
|
||||
|
||||
0.8.2
|
||||
=====
|
||||
|
||||
|
|
|
|||
|
|
@ -69,8 +69,6 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
public final static String PIG_INITIAL_ADDRESS = "PIG_INITIAL_ADDRESS";
|
||||
public final static String PIG_PARTITIONER = "PIG_PARTITIONER";
|
||||
|
||||
private static String UDFCONTEXT_SCHEMA_KEY_PREFIX = "cassandra.schema";
|
||||
|
||||
private final static ByteBuffer BOUND = ByteBufferUtil.EMPTY_BYTE_BUFFER;
|
||||
private static final Log logger = LogFactory.getLog(CassandraStorage.class);
|
||||
|
||||
|
|
@ -79,6 +77,8 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
private boolean slice_reverse = false;
|
||||
private String keyspace;
|
||||
private String column_family;
|
||||
private String loadSignature;
|
||||
private String storeSignature;
|
||||
|
||||
private Configuration conf;
|
||||
private RecordReader reader;
|
||||
|
|
@ -113,7 +113,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
if (!reader.nextKeyValue())
|
||||
return null;
|
||||
|
||||
CfDef cfDef = getCfDef();
|
||||
CfDef cfDef = getCfDef(loadSignature);
|
||||
ByteBuffer key = (ByteBuffer)reader.getCurrentKey();
|
||||
SortedMap<ByteBuffer,IColumn> cf = (SortedMap<ByteBuffer,IColumn>)reader.getCurrentValue();
|
||||
assert key != null && cf != null;
|
||||
|
|
@ -166,11 +166,11 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
return pair;
|
||||
}
|
||||
|
||||
private CfDef getCfDef()
|
||||
private CfDef getCfDef(String signature)
|
||||
{
|
||||
UDFContext context = UDFContext.getUDFContext();
|
||||
Properties property = context.getUDFProperties(CassandraStorage.class);
|
||||
return cfdefFromString(property.getProperty(getSchemaContextKey()));
|
||||
return cfdefFromString(property.getProperty(signature));
|
||||
}
|
||||
|
||||
private List<AbstractType> getDefaultMarshallers(CfDef cfDef) throws IOException
|
||||
|
|
@ -290,7 +290,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
}
|
||||
ConfigHelper.setInputColumnFamily(conf, keyspace, column_family);
|
||||
setConnectionInformation();
|
||||
initSchema();
|
||||
initSchema(loadSignature);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -299,9 +299,16 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
return location;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setUDFContextSignature(String signature)
|
||||
{
|
||||
this.loadSignature = signature;
|
||||
}
|
||||
|
||||
/* StoreFunc methods */
|
||||
public void setStoreFuncUDFContextSignature(String signature)
|
||||
{
|
||||
this.storeSignature = signature;
|
||||
}
|
||||
|
||||
public String relToAbsPathForStoreLocation(String location, Path curDir) throws IOException
|
||||
|
|
@ -315,7 +322,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
setLocationFromUri(location);
|
||||
ConfigHelper.setOutputColumnFamily(conf, keyspace, column_family);
|
||||
setConnectionInformation();
|
||||
initSchema();
|
||||
initSchema(storeSignature);
|
||||
}
|
||||
|
||||
public OutputFormat getOutputFormat()
|
||||
|
|
@ -347,7 +354,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
ByteBuffer key = objToBB(t.get(0));
|
||||
DefaultDataBag pairs = (DefaultDataBag) t.get(1);
|
||||
ArrayList<Mutation> mutationList = new ArrayList<Mutation>();
|
||||
CfDef cfDef = getCfDef();
|
||||
CfDef cfDef = getCfDef(storeSignature);
|
||||
List<AbstractType> marshallers = getDefaultMarshallers(cfDef);
|
||||
Map<ByteBuffer,AbstractType> validators = getValidatorMap(cfDef);
|
||||
try
|
||||
|
|
@ -412,7 +419,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
}
|
||||
catch (ClassCastException e)
|
||||
{
|
||||
throw new IOException(e + " Output must be (key, {(column,value)...}) for ColumnFamily or (key, {supercolumn:{(column,value)...}...}) for SuperColumnFamily");
|
||||
throw new IOException(e + " Output must be (key, {(column,value)...}) for ColumnFamily or (key, {supercolumn:{(column,value)...}...}) for SuperColumnFamily", e);
|
||||
}
|
||||
try
|
||||
{
|
||||
|
|
@ -430,14 +437,13 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
|
||||
/* Methods to get the column family schema from Cassandra */
|
||||
|
||||
private void initSchema()
|
||||
private void initSchema(String signature)
|
||||
{
|
||||
UDFContext context = UDFContext.getUDFContext();
|
||||
Properties property = context.getUDFProperties(CassandraStorage.class);
|
||||
|
||||
String schemaContextKey = getSchemaContextKey();
|
||||
// Only get the schema if we haven't already gotten it
|
||||
if (!property.containsKey(schemaContextKey))
|
||||
if (!property.containsKey(signature))
|
||||
{
|
||||
Cassandra.Client client = null;
|
||||
try
|
||||
|
|
@ -455,7 +461,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
break;
|
||||
}
|
||||
}
|
||||
property.setProperty(schemaContextKey, cfdefToString(cfDef));
|
||||
property.setProperty(signature, cfdefToString(cfDef));
|
||||
}
|
||||
catch (TException e)
|
||||
{
|
||||
|
|
@ -521,14 +527,4 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface
|
|||
}
|
||||
return cfDef;
|
||||
}
|
||||
|
||||
private String getSchemaContextKey()
|
||||
{
|
||||
StringBuilder sb = new StringBuilder(UDFCONTEXT_SCHEMA_KEY_PREFIX);
|
||||
sb.append('.');
|
||||
sb.append(keyspace);
|
||||
sb.append('.');
|
||||
sb.append(column_family);
|
||||
return sb.toString();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
h1. Cassandra Query Language (CQL) v1.0.0
|
||||
h1. Cassandra Query Language (CQL) v1.1.0
|
||||
|
||||
h2. Table of Contents
|
||||
|
||||
|
|
@ -364,5 +364,8 @@ Versioning of the CQL language adheres to the "Semantic Versioning":http://semve
|
|||
h1. Changes
|
||||
|
||||
pre.
|
||||
Sat, 01 Jun 2011 15:58:00 -0600 - Pavel Yaskevich
|
||||
* Updated to support ALTER (CASSANDRA-1709)
|
||||
|
||||
Tue, 22 Mar 2011 18:10:28 -0700 - Eric Evans <eevans@rackspace.com>
|
||||
* Initial version, 1.0.0
|
||||
|
|
|
|||
|
|
@ -204,6 +204,9 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean
|
|||
}
|
||||
waited = 0;
|
||||
// then wait for the correct schema version.
|
||||
// usually we use DD.getDefsVersion, which checks the local schema uuid as stored in the system table.
|
||||
// here we check the one in gossip instead; this serves as a canary to warn us if we introduce a bug that
|
||||
// causes the two to diverge (see CASSANDRA-2946)
|
||||
while (!gossiper.getEndpointStateForEndpoint(endpoint).getApplicationState(ApplicationState.SCHEMA).value.equals(
|
||||
gossiper.getEndpointStateForEndpoint(FBUtilities.getBroadcastAddress()).getApplicationState(ApplicationState.SCHEMA).value))
|
||||
{
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ public enum ApplicationState
|
|||
DC,
|
||||
RACK,
|
||||
RELEASE_VERSION,
|
||||
REMOVAL_COORDINATOR,
|
||||
INTERNAL_IP,
|
||||
// pad to allow adding new states to existing cluster
|
||||
X1,
|
||||
|
|
|
|||
|
|
@ -27,6 +27,7 @@ import java.util.*;
|
|||
import java.util.Map.Entry;
|
||||
import java.util.concurrent.*;
|
||||
|
||||
import org.apache.cassandra.dht.Token;
|
||||
import org.apache.cassandra.db.SystemTable;
|
||||
import org.apache.cassandra.net.MessageProducer;
|
||||
import org.apache.cassandra.config.ConfigurationException;
|
||||
|
|
@ -58,6 +59,8 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
private static final RetryingScheduledThreadPoolExecutor executor = new RetryingScheduledThreadPoolExecutor("GossipTasks");
|
||||
|
||||
static final ApplicationState[] STATES = ApplicationState.values();
|
||||
static final List<String> DEAD_STATES = Arrays.asList(VersionedValue.REMOVING_TOKEN, VersionedValue.REMOVED_TOKEN, VersionedValue.STATUS_LEFT);
|
||||
|
||||
private ScheduledFuture<?> scheduledGossipTask;
|
||||
public final static int intervalInMillis = 1000;
|
||||
public final static int QUARANTINE_DELAY = StorageService.RING_DELAY * 2;
|
||||
|
|
@ -264,17 +267,21 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
}
|
||||
|
||||
/**
|
||||
* Removes the endpoint from unreachable endpoint set
|
||||
* Removes the endpoint from gossip completely
|
||||
*
|
||||
* @param endpoint endpoint to be removed from the current membership.
|
||||
*/
|
||||
private void evictFromMembership(InetAddress endpoint)
|
||||
{
|
||||
unreachableEndpoints.remove(endpoint);
|
||||
endpointStateMap.remove(endpoint);
|
||||
justRemovedEndpoints.put(endpoint, System.currentTimeMillis());
|
||||
if (logger.isDebugEnabled())
|
||||
logger.debug("evicting " + endpoint + " from gossip");
|
||||
}
|
||||
|
||||
/**
|
||||
* Removes the endpoint completely from Gossip
|
||||
* Removes the endpoint from Gossip but retains endpoint state
|
||||
*/
|
||||
public void removeEndpoint(InetAddress endpoint)
|
||||
{
|
||||
|
|
@ -288,6 +295,8 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
FailureDetector.instance.remove(endpoint);
|
||||
versions.remove(endpoint);
|
||||
justRemovedEndpoints.put(endpoint, System.currentTimeMillis());
|
||||
if (logger.isDebugEnabled())
|
||||
logger.debug("removing endpoint " + endpoint);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -328,6 +337,67 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* This method will begin removing an existing endpoint from the cluster by spoofing its state
|
||||
* This should never be called unless this coordinator has had 'removetoken' invoked
|
||||
*
|
||||
* @param endpoint - the endpoint being removed
|
||||
* @param token - the token being removed
|
||||
* @param mytoken - my own token for replication coordination
|
||||
*/
|
||||
public void advertiseRemoving(InetAddress endpoint, Token token, Token mytoken)
|
||||
{
|
||||
EndpointState epState = endpointStateMap.get(endpoint);
|
||||
// remember this node's generation
|
||||
int generation = epState.getHeartBeatState().getGeneration();
|
||||
logger.info("Removing token: " + token);
|
||||
logger.info("Sleeping for " + StorageService.RING_DELAY + "ms to ensure " + endpoint + " does not change");
|
||||
try
|
||||
{
|
||||
Thread.sleep(StorageService.RING_DELAY);
|
||||
}
|
||||
catch (InterruptedException e)
|
||||
{
|
||||
throw new AssertionError(e);
|
||||
}
|
||||
// make sure it did not change
|
||||
epState = endpointStateMap.get(endpoint);
|
||||
if (epState.getHeartBeatState().getGeneration() != generation)
|
||||
throw new RuntimeException("Endpoint " + endpoint + " generation changed while trying to remove it");
|
||||
// update the other node's generation to mimic it as if it had changed it itself
|
||||
logger.info("Advertising removal for " + endpoint);
|
||||
epState.updateTimestamp(); // make sure we don't evict it too soon
|
||||
epState.getHeartBeatState().forceNewerGenerationUnsafe();
|
||||
epState.addApplicationState(ApplicationState.STATUS, StorageService.instance.valueFactory.removingNonlocal(token));
|
||||
epState.addApplicationState(ApplicationState.REMOVAL_COORDINATOR, StorageService.instance.valueFactory.removalCoordinator(mytoken));
|
||||
endpointStateMap.put(endpoint, epState);
|
||||
}
|
||||
|
||||
/**
|
||||
* Handles switching the endpoint's state from REMOVING_TOKEN to REMOVED_TOKEN
|
||||
* This should only be called after advertiseRemoving
|
||||
* @param endpoint
|
||||
* @param token
|
||||
*/
|
||||
public void advertiseTokenRemoved(InetAddress endpoint, Token token)
|
||||
{
|
||||
EndpointState epState = endpointStateMap.get(endpoint);
|
||||
epState.updateTimestamp(); // make sure we don't evict it too soon
|
||||
epState.getHeartBeatState().forceNewerGenerationUnsafe();
|
||||
epState.addApplicationState(ApplicationState.STATUS, StorageService.instance.valueFactory.removedNonlocal(token));
|
||||
logger.info("Completing removal of " + endpoint);
|
||||
endpointStateMap.put(endpoint, epState);
|
||||
// ensure at least one gossip round occurs before returning
|
||||
try
|
||||
{
|
||||
Thread.sleep(intervalInMillis * 2);
|
||||
}
|
||||
catch (InterruptedException e)
|
||||
{
|
||||
throw new AssertionError(e);
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isKnownEndpoint(InetAddress endpoint)
|
||||
{
|
||||
return endpointStateMap.containsKey(endpoint);
|
||||
|
|
@ -456,23 +526,18 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
{
|
||||
long duration = now - epState.getUpdateTimestamp();
|
||||
|
||||
if (StorageService.instance.getTokenMetadata().isMember(endpoint))
|
||||
epState.setHasToken(true);
|
||||
// check if this is a fat client. fat clients are removed automatically from
|
||||
// gosip after FatClientTimeout
|
||||
if (!epState.hasToken() && !epState.isAlive() && (duration > FatClientTimeout))
|
||||
if (!epState.hasToken() && !epState.isAlive() && !justRemovedEndpoints.containsKey(endpoint) && (duration > FatClientTimeout))
|
||||
{
|
||||
if (StorageService.instance.getTokenMetadata().isMember(endpoint))
|
||||
epState.setHasToken(true);
|
||||
else
|
||||
{
|
||||
if (!justRemovedEndpoints.containsKey(endpoint)) // if the node was decommissioned, it will have been removed but still appear as a fat client
|
||||
{
|
||||
logger.info("FatClient " + endpoint + " has been silent for " + FatClientTimeout + "ms, removing from gossip");
|
||||
removeEndpoint(endpoint); // after quarantine justRemoveEndpoints will remove the state
|
||||
}
|
||||
}
|
||||
logger.info("FatClient " + endpoint + " has been silent for " + FatClientTimeout + "ms, removing from gossip");
|
||||
removeEndpoint(endpoint); // will put it in justRemovedEndpoints to respect quarantine delay
|
||||
evictFromMembership(endpoint); // can get rid of the state immediately
|
||||
}
|
||||
|
||||
if ( !epState.isAlive() && (duration > aVeryLongTime) )
|
||||
if ( !epState.isAlive() && (duration > aVeryLongTime) && (!StorageService.instance.getTokenMetadata().isMember(endpoint)))
|
||||
{
|
||||
evictFromMembership(endpoint);
|
||||
}
|
||||
|
|
@ -488,7 +553,6 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
if (logger.isDebugEnabled())
|
||||
logger.debug(QUARANTINE_DELAY + " elapsed, " + entry.getKey() + " gossip quarantine over");
|
||||
justRemovedEndpoints.remove(entry.getKey());
|
||||
endpointStateMap.remove(entry.getKey());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -585,6 +649,7 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
int remoteGeneration = remoteEndpointState.getHeartBeatState().getGeneration();
|
||||
if ( remoteGeneration > localGeneration )
|
||||
{
|
||||
localEndpointState.updateTimestamp();
|
||||
fd.report(endpoint);
|
||||
return;
|
||||
}
|
||||
|
|
@ -595,6 +660,7 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
int remoteVersion = remoteEndpointState.getHeartBeatState().getHeartBeatVersion();
|
||||
if ( remoteVersion > localVersion )
|
||||
{
|
||||
localEndpointState.updateTimestamp();
|
||||
fd.report(endpoint);
|
||||
}
|
||||
}
|
||||
|
|
@ -607,6 +673,7 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
if (logger.isTraceEnabled())
|
||||
logger.trace("marking as alive {}", addr);
|
||||
localState.markAlive();
|
||||
localState.updateTimestamp(); // prevents doStatusCheck from racing us and evicting if it was down > aVeryLongTime
|
||||
liveEndpoints.add(addr);
|
||||
unreachableEndpoints.remove(addr);
|
||||
logger.info("InetAddress {} is now UP", addr);
|
||||
|
|
@ -638,10 +705,13 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
*/
|
||||
private void handleMajorStateChange(InetAddress ep, EndpointState epState)
|
||||
{
|
||||
if (endpointStateMap.get(ep) != null)
|
||||
logger.info("Node {} has restarted, now UP again", ep);
|
||||
else
|
||||
logger.info("Node {} is now part of the cluster", ep);
|
||||
if (epState.getApplicationState(ApplicationState.STATUS) != null && !isDeadState(epState.getApplicationState(ApplicationState.STATUS).value))
|
||||
{
|
||||
if (endpointStateMap.get(ep) != null)
|
||||
logger.info("Node {} has restarted, now UP again", ep);
|
||||
else
|
||||
logger.info("Node {} is now part of the cluster", ep);
|
||||
}
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Adding endpoint state for " + ep);
|
||||
endpointStateMap.put(ep, epState);
|
||||
|
|
@ -651,11 +721,31 @@ public class Gossiper implements IFailureDetectionEventListener
|
|||
for (IEndpointStateChangeSubscriber subscriber : subscribers)
|
||||
subscriber.onDead(ep, epState);
|
||||
}
|
||||
markAlive(ep, epState);
|
||||
if (epState.getApplicationState(ApplicationState.STATUS) != null && !isDeadState(epState.getApplicationState(ApplicationState.STATUS).value))
|
||||
markAlive(ep, epState);
|
||||
else
|
||||
{
|
||||
logger.debug("Not marking " + ep + " alive due to dead state");
|
||||
epState.markDead();
|
||||
epState.setHasToken(true); // fat clients won't have a dead state
|
||||
}
|
||||
for (IEndpointStateChangeSubscriber subscriber : subscribers)
|
||||
subscriber.onJoin(ep, epState);
|
||||
}
|
||||
|
||||
private Boolean isDeadState(String value)
|
||||
{
|
||||
String[] pieces = value.split(VersionedValue.DELIMITER_STR, -1);
|
||||
assert (pieces.length > 0);
|
||||
String state = pieces[0];
|
||||
for (String deadstate : DEAD_STATES)
|
||||
{
|
||||
if (state.equals(deadstate))
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
void applyStateLocally(Map<InetAddress, EndpointState> epStateMap)
|
||||
{
|
||||
for (Entry<InetAddress, EndpointState> entry : epStateMap.entrySet())
|
||||
|
|
|
|||
|
|
@ -71,6 +71,11 @@ class HeartBeatState
|
|||
{
|
||||
return version;
|
||||
}
|
||||
|
||||
void forceNewerGenerationUnsafe()
|
||||
{
|
||||
generation += 1;
|
||||
}
|
||||
}
|
||||
|
||||
class HeartBeatStateSerializer implements ICompactSerializer<HeartBeatState>
|
||||
|
|
|
|||
|
|
@ -49,7 +49,7 @@ public class VersionedValue implements Comparable<VersionedValue>
|
|||
public final static char DELIMITER = ',';
|
||||
public final static String DELIMITER_STR = new String(new char[] { DELIMITER });
|
||||
|
||||
// values for State.STATUS
|
||||
// values for ApplicationState.STATUS
|
||||
public final static String STATUS_BOOTSTRAPPING = "BOOT";
|
||||
public final static String STATUS_NORMAL = "NORMAL";
|
||||
public final static String STATUS_LEAVING = "LEAVING";
|
||||
|
|
@ -59,6 +59,9 @@ public class VersionedValue implements Comparable<VersionedValue>
|
|||
public final static String REMOVING_TOKEN = "removing";
|
||||
public final static String REMOVED_TOKEN = "removed";
|
||||
|
||||
// values for ApplicationState.REMOVAL_COORDINATOR
|
||||
public final static String REMOVAL_COORDINATOR = "REMOVER";
|
||||
|
||||
public final int version;
|
||||
public final String value;
|
||||
|
||||
|
|
@ -129,20 +132,19 @@ public class VersionedValue implements Comparable<VersionedValue>
|
|||
return new VersionedValue(VersionedValue.STATUS_MOVING + VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(token));
|
||||
}
|
||||
|
||||
public VersionedValue removingNonlocal(Token localToken, Token token)
|
||||
public VersionedValue removingNonlocal(Token token)
|
||||
{
|
||||
return new VersionedValue(VersionedValue.STATUS_NORMAL
|
||||
+ VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(localToken)
|
||||
+ VersionedValue.DELIMITER + VersionedValue.REMOVING_TOKEN
|
||||
+ VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(token));
|
||||
return new VersionedValue(VersionedValue.REMOVING_TOKEN + VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(token));
|
||||
}
|
||||
|
||||
public VersionedValue removedNonlocal(Token localToken, Token token)
|
||||
public VersionedValue removedNonlocal(Token token)
|
||||
{
|
||||
return new VersionedValue(VersionedValue.STATUS_NORMAL
|
||||
+ VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(localToken)
|
||||
+ VersionedValue.DELIMITER + VersionedValue.REMOVED_TOKEN
|
||||
+ VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(token));
|
||||
return new VersionedValue(VersionedValue.REMOVED_TOKEN + VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(token));
|
||||
}
|
||||
|
||||
public VersionedValue removalCoordinator(Token token)
|
||||
{
|
||||
return new VersionedValue(VersionedValue.REMOVAL_COORDINATOR + VersionedValue.DELIMITER + partitioner.getTokenFactory().toString(token));
|
||||
}
|
||||
|
||||
public VersionedValue datacenter(String dcId)
|
||||
|
|
|
|||
|
|
@ -472,6 +472,10 @@ public final class MessagingService implements MessagingServiceMBean
|
|||
|
||||
public void receive(Message message, String id)
|
||||
{
|
||||
if (logger_.isTraceEnabled())
|
||||
logger_.trace(FBUtilities.getLocalAddress() + " received " + message.getVerb()
|
||||
+ " from " + id + "@" + message.getFrom());
|
||||
|
||||
message = SinkManager.processServerMessage(message, id);
|
||||
if (message == null)
|
||||
return;
|
||||
|
|
|
|||
|
|
@ -119,6 +119,7 @@ public abstract class AbstractCassandraDaemon implements CassandraDaemon
|
|||
{
|
||||
logger.info("JVM vendor/version: {}/{}", System.getProperty("java.vm.name"), System.getProperty("java.version") );
|
||||
logger.info("Heap size: {}/{}", Runtime.getRuntime().totalMemory(), Runtime.getRuntime().maxMemory());
|
||||
logger.info("Classpath: {}", System.getProperty("java.class.path"));
|
||||
CLibrary.tryMlockall();
|
||||
|
||||
listenPort = DatabaseDescriptor.getRpcPort();
|
||||
|
|
|
|||
|
|
@ -349,7 +349,7 @@ public class StorageProxy implements StorageProxyMBean
|
|||
{
|
||||
if (logger.isDebugEnabled())
|
||||
logger.debug("insert writing local " + rm.toString(true));
|
||||
Runnable runnable = new WrappedRunnable()
|
||||
Runnable runnable = new DroppableRunnable(StorageService.Verb.MUTATION)
|
||||
{
|
||||
public void runMayThrow() throws IOException
|
||||
{
|
||||
|
|
@ -431,7 +431,7 @@ public class StorageProxy implements StorageProxyMBean
|
|||
if (logger.isDebugEnabled())
|
||||
logger.debug("insert writing local & replicate " + mutation.toString(true));
|
||||
|
||||
Runnable runnable = new WrappedRunnable()
|
||||
Runnable runnable = new DroppableRunnable(StorageService.Verb.MUTATION)
|
||||
{
|
||||
public void runMayThrow() throws IOException
|
||||
{
|
||||
|
|
@ -447,7 +447,7 @@ public class StorageProxy implements StorageProxyMBean
|
|||
{
|
||||
// We do the replication on another stage because it involves a read (see CM.makeReplicationMutation)
|
||||
// and we want to avoid blocking too much the MUTATION stage
|
||||
StageManager.getStage(Stage.REPLICATE_ON_WRITE).execute(new WrappedRunnable()
|
||||
StageManager.getStage(Stage.REPLICATE_ON_WRITE).execute(new DroppableRunnable(StorageService.Verb.READ)
|
||||
{
|
||||
public void runMayThrow() throws IOException
|
||||
{
|
||||
|
|
@ -616,7 +616,7 @@ public class StorageProxy implements StorageProxyMBean
|
|||
return rows;
|
||||
}
|
||||
|
||||
static class LocalReadRunnable extends WrappedRunnable
|
||||
static class LocalReadRunnable extends DroppableRunnable
|
||||
{
|
||||
private final ReadCommand command;
|
||||
private final ReadCallback<Row> handler;
|
||||
|
|
@ -624,6 +624,7 @@ public class StorageProxy implements StorageProxyMBean
|
|||
|
||||
LocalReadRunnable(ReadCommand command, ReadCallback<Row> handler)
|
||||
{
|
||||
super(StorageService.Verb.READ);
|
||||
this.command = command;
|
||||
this.handler = handler;
|
||||
}
|
||||
|
|
@ -1078,4 +1079,35 @@ public class StorageProxy implements StorageProxyMBean
|
|||
{
|
||||
public void apply(IMutation mutation, Multimap<InetAddress, InetAddress> hintedEndpoints, IWriteResponseHandler responseHandler, String localDataCenter, ConsistencyLevel consistency_level) throws IOException;
|
||||
}
|
||||
|
||||
private static abstract class DroppableRunnable implements Runnable
|
||||
{
|
||||
private final long constructionTime = System.currentTimeMillis();
|
||||
private final StorageService.Verb verb;
|
||||
|
||||
public DroppableRunnable(StorageService.Verb verb)
|
||||
{
|
||||
this.verb = verb;
|
||||
}
|
||||
|
||||
public final void run()
|
||||
{
|
||||
if (System.currentTimeMillis() > constructionTime + DatabaseDescriptor.getRpcTimeout())
|
||||
{
|
||||
MessagingService.instance().incrementDroppedMessages(verb);
|
||||
return;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
runMayThrow();
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
abstract protected void runMayThrow() throws Exception;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -334,6 +334,11 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
}
|
||||
|
||||
public synchronized void initClient() throws IOException, ConfigurationException
|
||||
{
|
||||
initClient(RING_DELAY);
|
||||
}
|
||||
|
||||
public synchronized void initClient(int delay) throws IOException, ConfigurationException
|
||||
{
|
||||
if (initialized)
|
||||
{
|
||||
|
|
@ -352,7 +357,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
// sleep a while to allow gossip to warm up (the other nodes need to know about this one before they can reply).
|
||||
try
|
||||
{
|
||||
Thread.sleep(RING_DELAY);
|
||||
Thread.sleep(delay);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
|
|
@ -622,29 +627,35 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
}
|
||||
|
||||
/*
|
||||
* onChange only ever sees one ApplicationState piece change at a time, so we perform a kind of state machine here.
|
||||
* We are concerned with two events: knowing the token associated with an endpoint, and knowing its operation mode.
|
||||
* Nodes can start in either bootstrap or normal mode, and from bootstrap mode can change mode to normal.
|
||||
* A node in bootstrap mode needs to have pendingranges set in TokenMetadata; a node in normal mode
|
||||
* should instead be part of the token ring.
|
||||
* Handle the reception of a new particular ApplicationState for a particular endpoint. Note that the value of the
|
||||
* ApplicationState has not necessarily "changed" since the last known value, if we already received the same update
|
||||
* from somewhere else.
|
||||
*
|
||||
* onChange only ever sees one ApplicationState piece change at a time (even if many ApplicationState updates were
|
||||
* received at the same time), so we perform a kind of state machine here. We are concerned with two events: knowing
|
||||
* the token associated with an endpoint, and knowing its operation mode. Nodes can start in either bootstrap or
|
||||
* normal mode, and from bootstrap mode can change mode to normal. A node in bootstrap mode needs to have
|
||||
* pendingranges set in TokenMetadata; a node in normal mode should instead be part of the token ring.
|
||||
*
|
||||
* Normal MOVE_STATE progression of a node should be like this:
|
||||
* STATE_BOOTSTRAPPING,token
|
||||
* Normal progression of ApplicationState.STATUS values for a node should be like this:
|
||||
* STATUS_BOOTSTRAPPING,token
|
||||
* if bootstrapping. stays this way until all files are received.
|
||||
* STATE_NORMAL,token
|
||||
* STATUS_NORMAL,token
|
||||
* ready to serve reads and writes.
|
||||
* STATE_NORMAL,token,REMOVE_TOKEN,token
|
||||
* specialized normal state in which this node acts as a proxy to tell the cluster about a dead node whose
|
||||
* token is being removed. this value becomes the permanent state of this node (unless it coordinates another
|
||||
* removetoken in the future).
|
||||
* STATE_LEAVING,token
|
||||
* get ready to leave the cluster as part of a decommission or move
|
||||
* STATE_LEFT,token
|
||||
* set after decommission or move is completed.
|
||||
* STATE_MOVE,token
|
||||
* set if node if currently moving to a new token in the ring
|
||||
*
|
||||
* Note: Any time a node state changes from STATE_NORMAL, it will not be visible to new nodes. So it follows that
|
||||
* STATUS_LEAVING,token
|
||||
* get ready to leave the cluster as part of a decommission
|
||||
* STATUS_LEFT,token
|
||||
* set after decommission is completed.
|
||||
*
|
||||
* Other STATUS values that may be seen (possibly anywhere in the normal progression):
|
||||
* STATUS_MOVING,newtoken
|
||||
* set if node is currently moving to a new token in the ring
|
||||
* REMOVING_TOKEN,deadtoken
|
||||
* set if the node is dead and is being removed by its REMOVAL_COORDINATOR
|
||||
* REMOVED_TOKEN,deadtoken
|
||||
* set if the node is dead and has been removed by its REMOVAL_COORDINATOR
|
||||
*
|
||||
* Note: Any time a node state changes from STATUS_NORMAL, it will not be visible to new nodes. So it follows that
|
||||
* you should never bootstrap a new node during a removetoken, decommission or move.
|
||||
*/
|
||||
public void onChange(InetAddress endpoint, ApplicationState state, VersionedValue value)
|
||||
|
|
@ -665,6 +676,8 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
handleStateBootstrap(endpoint, pieces);
|
||||
else if (moveName.equals(VersionedValue.STATUS_NORMAL))
|
||||
handleStateNormal(endpoint, pieces);
|
||||
else if (moveName.equals(VersionedValue.REMOVING_TOKEN) || moveName.equals(VersionedValue.REMOVED_TOKEN))
|
||||
handleStateRemoving(endpoint, pieces);
|
||||
else if (moveName.equals(VersionedValue.STATUS_LEAVING))
|
||||
handleStateLeaving(endpoint, pieces);
|
||||
else if (moveName.equals(VersionedValue.STATUS_LEFT))
|
||||
|
|
@ -731,7 +744,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
* in reads.
|
||||
*
|
||||
* @param endpoint node
|
||||
* @param pieces STATE_NORMAL,token[,other_state,token]
|
||||
* @param pieces STATE_NORMAL,token
|
||||
*/
|
||||
private void handleStateNormal(InetAddress endpoint, String[] pieces)
|
||||
{
|
||||
|
|
@ -773,12 +786,6 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
endpoint, currentOwner, token, endpoint));
|
||||
}
|
||||
|
||||
if (pieces.length > 2)
|
||||
{
|
||||
assert pieces.length == 4;
|
||||
handleStateRemoving(endpoint, getPartitioner().getTokenFactory().fromString(pieces[3]), pieces[2]);
|
||||
}
|
||||
|
||||
if (tokenMetadata_.isMoving(endpoint)) // if endpoint was moving to a new token
|
||||
tokenMetadata_.removeFromMoving(endpoint);
|
||||
|
||||
|
|
@ -860,37 +867,50 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
* Handle notification that a node being actively removed from the ring via 'removetoken'
|
||||
*
|
||||
* @param endpoint node
|
||||
* @param state either REMOVED_TOKEN (node is gone) or REMOVING_TOKEN (replicas need to be restored)
|
||||
* @param pieces either REMOVED_TOKEN (node is gone) or REMOVING_TOKEN (replicas need to be restored)
|
||||
*/
|
||||
private void handleStateRemoving(InetAddress endpoint, Token removeToken, String state)
|
||||
private void handleStateRemoving(InetAddress endpoint, String[] pieces)
|
||||
{
|
||||
InetAddress removeEndpoint = tokenMetadata_.getEndpoint(removeToken);
|
||||
|
||||
if (removeEndpoint == null)
|
||||
return;
|
||||
|
||||
if (removeEndpoint.equals(FBUtilities.getBroadcastAddress()))
|
||||
String state = pieces[0];
|
||||
assert (pieces.length > 0);
|
||||
|
||||
if (endpoint.equals(FBUtilities.getBroadcastAddress()))
|
||||
{
|
||||
logger_.info("Received removeToken gossip about myself. Is this node a replacement for a removed one?");
|
||||
logger_.info("Received removeToken gossip about myself. Is this node rejoining after an explicit removetoken?");
|
||||
try
|
||||
{
|
||||
drain();
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (VersionedValue.REMOVED_TOKEN.equals(state))
|
||||
if (tokenMetadata_.isMember(endpoint))
|
||||
{
|
||||
excise(removeToken, removeEndpoint);
|
||||
}
|
||||
else if (VersionedValue.REMOVING_TOKEN.equals(state))
|
||||
{
|
||||
if (logger_.isDebugEnabled())
|
||||
logger_.debug("Token " + removeToken + " removed manually (endpoint was " + removeEndpoint + ")");
|
||||
Token removeToken = tokenMetadata_.getToken(endpoint);
|
||||
|
||||
// Note that the endpoint is being removed
|
||||
tokenMetadata_.addLeavingEndpoint(removeEndpoint);
|
||||
calculatePendingRanges();
|
||||
if (VersionedValue.REMOVED_TOKEN.equals(state))
|
||||
{
|
||||
excise(removeToken, endpoint);
|
||||
}
|
||||
else if (VersionedValue.REMOVING_TOKEN.equals(state))
|
||||
{
|
||||
if (logger_.isDebugEnabled())
|
||||
logger_.debug("Token " + removeToken + " removed manually (endpoint was " + endpoint + ")");
|
||||
|
||||
// grab any data we are now responsible for and notify responsible node
|
||||
restoreReplicaCount(removeEndpoint, endpoint);
|
||||
}
|
||||
// Note that the endpoint is being removed
|
||||
tokenMetadata_.addLeavingEndpoint(endpoint);
|
||||
calculatePendingRanges();
|
||||
|
||||
// find the endpoint coordinating this removal that we need to notify when we're done
|
||||
String[] coordinator = Gossiper.instance.getEndpointStateForEndpoint(endpoint).getApplicationState(ApplicationState.REMOVAL_COORDINATOR).value.split(VersionedValue.DELIMITER_STR, -1);
|
||||
Token coordtoken = getPartitioner().getTokenFactory().fromString(coordinator[1]);
|
||||
// grab any data we are now responsible for and notify responsible node
|
||||
restoreReplicaCount(endpoint, tokenMetadata_.getEndpoint(coordtoken));
|
||||
}
|
||||
} // not a member, nothing to do
|
||||
}
|
||||
|
||||
private void excise(Token token, InetAddress endpoint)
|
||||
|
|
@ -1059,6 +1079,8 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
// notify the remote token
|
||||
Message msg = new Message(local, StorageService.Verb.REPLICATION_FINISHED, new byte[0], Gossiper.instance.getVersion(remote));
|
||||
IFailureDetector failureDetector = FailureDetector.instance;
|
||||
if (logger_.isDebugEnabled())
|
||||
logger_.debug("Notifying " + remote.toString() + " of replication completion\n");
|
||||
while (failureDetector.isAlive(remote))
|
||||
{
|
||||
IAsyncResult iar = MessagingService.instance().sendRR(msg, remote);
|
||||
|
|
@ -1993,9 +2015,14 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
*/
|
||||
public void forceRemoveCompletion()
|
||||
{
|
||||
if (!replicatingNodes.isEmpty())
|
||||
if (!replicatingNodes.isEmpty() || !tokenMetadata_.getLeavingEndpoints().isEmpty())
|
||||
{
|
||||
logger_.warn("Removal not confirmed for for " + StringUtils.join(this.replicatingNodes, ","));
|
||||
for (InetAddress endpoint : tokenMetadata_.getLeavingEndpoints())
|
||||
{
|
||||
Gossiper.instance.advertiseTokenRemoved(endpoint, tokenMetadata_.getToken(endpoint));
|
||||
tokenMetadata_.removeEndpoint(endpoint);
|
||||
}
|
||||
replicatingNodes.clear();
|
||||
}
|
||||
else
|
||||
|
|
@ -2059,9 +2086,9 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
|
||||
tokenMetadata_.addLeavingEndpoint(endpoint);
|
||||
calculatePendingRanges();
|
||||
// bundle two states together. include this nodes state to keep the status quo,
|
||||
// but indicate the leaving token so that it can be dealt with.
|
||||
Gossiper.instance.addLocalApplicationState(ApplicationState.STATUS, valueFactory.removingNonlocal(localToken, token));
|
||||
// the gossiper will handle spoofing this node's state to REMOVING_TOKEN for us
|
||||
// we add our own token so other nodes to let us know when they're done
|
||||
Gossiper.instance.advertiseRemoving(endpoint, token, localToken);
|
||||
|
||||
// kick off streaming commands
|
||||
restoreReplicaCount(endpoint, myAddress);
|
||||
|
|
@ -2081,8 +2108,8 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
|
||||
excise(token, endpoint);
|
||||
|
||||
// indicate the token has left
|
||||
Gossiper.instance.addLocalApplicationState(ApplicationState.STATUS, valueFactory.removedNonlocal(localToken, token));
|
||||
// gossiper will indicate the token has left
|
||||
Gossiper.instance.advertiseTokenRemoved(endpoint, token);
|
||||
|
||||
replicatingNodes.clear();
|
||||
removingNode = null;
|
||||
|
|
@ -2090,8 +2117,18 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe
|
|||
|
||||
public void confirmReplication(InetAddress node)
|
||||
{
|
||||
assert !replicatingNodes.isEmpty();
|
||||
replicatingNodes.remove(node);
|
||||
// replicatingNodes can be empty in the case where this node used to be a removal coordinator,
|
||||
// but restarted before all 'replication finished' messages arrived. In that case, we'll
|
||||
// still go ahead and acknowledge it.
|
||||
if (!replicatingNodes.isEmpty())
|
||||
{
|
||||
replicatingNodes.remove(node);
|
||||
}
|
||||
else
|
||||
{
|
||||
logger_.info("Received unexpected REPLICATION_FINISHED message from " + node
|
||||
+ ". Was this node recently a removal coordinator?");
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isClientMode()
|
||||
|
|
|
|||
|
|
@ -30,6 +30,6 @@ public class InitClientTest // extends CleanupHelper
|
|||
@Test
|
||||
public void testInitClientStartup() throws IOException, ConfigurationException
|
||||
{
|
||||
StorageService.instance.initClient();
|
||||
StorageService.instance.initClient(0);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -36,7 +36,7 @@ public class StorageServiceClientTest
|
|||
{
|
||||
CleanupHelper.mkdirs();
|
||||
CleanupHelper.cleanup();
|
||||
StorageService.instance.initClient();
|
||||
StorageService.instance.initClient(0);
|
||||
|
||||
// verify that no storage directories were created.
|
||||
for (String path : DatabaseDescriptor.getAllDataFileLocations())
|
||||
|
|
|
|||
Loading…
Reference in New Issue