Merge branch 'cassandra-3.X' into trunk

* cassandra-3.X:
  Use new token allocation for non bootstrap case as well.
This commit is contained in:
Dikang Gu 2017-01-05 12:11:58 -08:00
commit d384e781d6
5 changed files with 48 additions and 45 deletions

View File

@ -10,6 +10,7 @@
3.12 3.12
* Use new token allocation for non bootstrap case as well (CASSANDRA-13080)
* Avoid byte-array copy when key cache is disabled (CASSANDRA-13084) * Avoid byte-array copy when key cache is disabled (CASSANDRA-13084)
* More fixes to the TokenAllocator (CASSANDRA-12990) * More fixes to the TokenAllocator (CASSANDRA-12990)
* Require forceful decommission if number of nodes is less than replication factor (CASSANDRA-12510) * Require forceful decommission if number of nodes is less than replication factor (CASSANDRA-12510)

View File

@ -33,12 +33,15 @@ import org.apache.cassandra.db.TypeSizes;
import org.apache.cassandra.dht.tokenallocator.TokenAllocation; import org.apache.cassandra.dht.tokenallocator.TokenAllocation;
import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.gms.FailureDetector; import org.apache.cassandra.gms.FailureDetector;
import org.apache.cassandra.gms.Gossiper;
import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.IVersionedSerializer;
import org.apache.cassandra.io.util.DataInputPlus; import org.apache.cassandra.io.util.DataInputPlus;
import org.apache.cassandra.io.util.DataOutputPlus; import org.apache.cassandra.io.util.DataOutputPlus;
import org.apache.cassandra.locator.AbstractReplicationStrategy; import org.apache.cassandra.locator.AbstractReplicationStrategy;
import org.apache.cassandra.locator.TokenMetadata; import org.apache.cassandra.locator.TokenMetadata;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.streaming.*; import org.apache.cassandra.streaming.*;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.progress.ProgressEvent; import org.apache.cassandra.utils.progress.ProgressEvent;
import org.apache.cassandra.utils.progress.ProgressEventNotifierSupport; import org.apache.cassandra.utils.progress.ProgressEventNotifierSupport;
import org.apache.cassandra.utils.progress.ProgressEventType; import org.apache.cassandra.utils.progress.ProgressEventType;
@ -155,7 +158,7 @@ public class BootStrapper extends ProgressEventNotifierSupport
* otherwise, if allocationKeyspace is specified use the token allocation algorithm to generate suitable tokens * otherwise, if allocationKeyspace is specified use the token allocation algorithm to generate suitable tokens
* else choose num_tokens tokens at random * else choose num_tokens tokens at random
*/ */
public static Collection<Token> getBootstrapTokens(final TokenMetadata metadata, InetAddress address) throws ConfigurationException public static Collection<Token> getBootstrapTokens(final TokenMetadata metadata, InetAddress address, int schemaWaitDelay) throws ConfigurationException
{ {
String allocationKeyspace = DatabaseDescriptor.getAllocateTokensForKeyspace(); String allocationKeyspace = DatabaseDescriptor.getAllocateTokensForKeyspace();
Collection<String> initialTokens = DatabaseDescriptor.getInitialTokens(); Collection<String> initialTokens = DatabaseDescriptor.getInitialTokens();
@ -171,7 +174,7 @@ public class BootStrapper extends ProgressEventNotifierSupport
throw new ConfigurationException("num_tokens must be >= 1"); throw new ConfigurationException("num_tokens must be >= 1");
if (allocationKeyspace != null) if (allocationKeyspace != null)
return allocateTokens(metadata, address, allocationKeyspace, numTokens); return allocateTokens(metadata, address, allocationKeyspace, numTokens, schemaWaitDelay);
if (numTokens == 1) if (numTokens == 1)
logger.warn("Picking random token for a single vnode. You should probably add more vnodes and/or use the automatic token allocation mechanism."); logger.warn("Picking random token for a single vnode. You should probably add more vnodes and/or use the automatic token allocation mechanism.");
@ -182,7 +185,7 @@ public class BootStrapper extends ProgressEventNotifierSupport
private static Collection<Token> getSpecifiedTokens(final TokenMetadata metadata, private static Collection<Token> getSpecifiedTokens(final TokenMetadata metadata,
Collection<String> initialTokens) Collection<String> initialTokens)
{ {
logger.trace("tokens manually specified as {}", initialTokens); logger.info("tokens manually specified as {}", initialTokens);
List<Token> tokens = new ArrayList<>(initialTokens.size()); List<Token> tokens = new ArrayList<>(initialTokens.size());
for (String tokenString : initialTokens) for (String tokenString : initialTokens)
{ {
@ -197,8 +200,13 @@ public class BootStrapper extends ProgressEventNotifierSupport
static Collection<Token> allocateTokens(final TokenMetadata metadata, static Collection<Token> allocateTokens(final TokenMetadata metadata,
InetAddress address, InetAddress address,
String allocationKeyspace, String allocationKeyspace,
int numTokens) int numTokens,
int schemaWaitDelay)
{ {
StorageService.instance.waitForSchema(schemaWaitDelay);
if (!FBUtilities.getBroadcastAddress().equals(InetAddress.getLoopbackAddress()))
Gossiper.waitToSettle();
Keyspace ks = Keyspace.open(allocationKeyspace); Keyspace ks = Keyspace.open(allocationKeyspace);
if (ks == null) if (ks == null)
throw new ConfigurationException("Problem opening token allocation keyspace " + allocationKeyspace); throw new ConfigurationException("Problem opening token allocation keyspace " + allocationKeyspace);
@ -216,6 +224,8 @@ public class BootStrapper extends ProgressEventNotifierSupport
if (metadata.getEndpoint(token) == null) if (metadata.getEndpoint(token) == null)
tokens.add(token); tokens.add(token);
} }
logger.info("Generated random tokens. tokens are {}", tokens);
return tokens; return tokens;
} }

View File

@ -53,9 +53,6 @@ public class TokenAllocation
final InetAddress endpoint, final InetAddress endpoint,
int numTokens) int numTokens)
{ {
if (!FBUtilities.getBroadcastAddress().equals(InetAddress.getLoopbackAddress()))
Gossiper.waitToSettle();
TokenMetadata tokenMetadataCopy = tokenMetadata.cloneOnlyTokenMap(); TokenMetadata tokenMetadataCopy = tokenMetadata.cloneOnlyTokenMap();
StrategyAdapter strategy = getStrategy(tokenMetadataCopy, rs, endpoint); StrategyAdapter strategy = getStrategy(tokenMetadataCopy, rs, endpoint);
Collection<Token> tokens = create(tokenMetadata, strategy).addUnit(endpoint, numTokens); Collection<Token> tokens = create(tokenMetadata, strategy).addUnit(endpoint, numTokens);

View File

@ -709,7 +709,12 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
private boolean shouldBootstrap() private boolean shouldBootstrap()
{ {
return DatabaseDescriptor.isAutoBootstrap() && !SystemKeyspace.bootstrapComplete() && !DatabaseDescriptor.getSeeds().contains(FBUtilities.getBroadcastAddress()); return DatabaseDescriptor.isAutoBootstrap() && !SystemKeyspace.bootstrapComplete() && !isSeed();
}
public static boolean isSeed()
{
return DatabaseDescriptor.getSeeds().contains(FBUtilities.getBroadcastAddress());
} }
private void prepareToJoin() throws ConfigurationException private void prepareToJoin() throws ConfigurationException
@ -792,6 +797,29 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
} }
} }
public void waitForSchema(int delay)
{
// first sleep the delay to make sure we see all our peers
for (int i = 0; i < delay; i += 1000)
{
// if we see schema, we can proceed to the next check directly
if (!Schema.instance.getVersion().equals(SchemaConstants.emptyVersion))
{
logger.debug("got schema: {}", Schema.instance.getVersion());
break;
}
Uninterruptibles.sleepUninterruptibly(1, TimeUnit.SECONDS);
}
// if our schema hasn't matched yet, wait until it has
// we do this by waiting for all in-flight migration requests and responses to complete
// (post CASSANDRA-1391 we don't expect this to be necessary very often, but it doesn't hurt to be careful)
if (!MigrationManager.isReadyForBootstrap())
{
setMode(Mode.JOINING, "waiting for schema information to complete", true);
MigrationManager.waitUntilReadyForBootstrap();
}
}
private void joinTokenRing(int delay) throws ConfigurationException private void joinTokenRing(int delay) throws ConfigurationException
{ {
joined = true; joined = true;
@ -828,25 +856,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
else else
SystemKeyspace.setBootstrapState(SystemKeyspace.BootstrapState.IN_PROGRESS); SystemKeyspace.setBootstrapState(SystemKeyspace.BootstrapState.IN_PROGRESS);
setMode(Mode.JOINING, "waiting for ring information", true); setMode(Mode.JOINING, "waiting for ring information", true);
// first sleep the delay to make sure we see all our peers waitForSchema(delay);
for (int i = 0; i < delay; i += 1000)
{
// if we see schema, we can proceed to the next check directly
if (!Schema.instance.getVersion().equals(SchemaConstants.emptyVersion))
{
logger.debug("got schema: {}", Schema.instance.getVersion());
break;
}
Uninterruptibles.sleepUninterruptibly(1, TimeUnit.SECONDS);
}
// if our schema hasn't matched yet, wait until it has
// we do this by waiting for all in-flight migration requests and responses to complete
// (post CASSANDRA-1391 we don't expect this to be necessary very often, but it doesn't hurt to be careful)
if (!MigrationManager.isReadyForBootstrap())
{
setMode(Mode.JOINING, "waiting for schema information to complete", true);
MigrationManager.waitUntilReadyForBootstrap();
}
setMode(Mode.JOINING, "schema complete, ready to bootstrap", true); setMode(Mode.JOINING, "schema complete, ready to bootstrap", true);
setMode(Mode.JOINING, "waiting for pending range calculation", true); setMode(Mode.JOINING, "waiting for pending range calculation", true);
PendingRangeCalculatorService.instance.blockUntilFinished(); PendingRangeCalculatorService.instance.blockUntilFinished();
@ -873,7 +883,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
throw new UnsupportedOperationException(s); throw new UnsupportedOperationException(s);
} }
setMode(Mode.JOINING, "getting bootstrap token", true); setMode(Mode.JOINING, "getting bootstrap token", true);
bootstrapTokens = BootStrapper.getBootstrapTokens(tokenMetadata, FBUtilities.getBroadcastAddress()); bootstrapTokens = BootStrapper.getBootstrapTokens(tokenMetadata, FBUtilities.getBroadcastAddress(), delay);
} }
else else
{ {
@ -929,22 +939,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
bootstrapTokens = SystemKeyspace.getSavedTokens(); bootstrapTokens = SystemKeyspace.getSavedTokens();
if (bootstrapTokens.isEmpty()) if (bootstrapTokens.isEmpty())
{ {
Collection<String> initialTokens = DatabaseDescriptor.getInitialTokens(); bootstrapTokens = BootStrapper.getBootstrapTokens(tokenMetadata, FBUtilities.getBroadcastAddress(), delay);
if (initialTokens.size() < 1)
{
bootstrapTokens = BootStrapper.getRandomTokens(tokenMetadata, DatabaseDescriptor.getNumTokens());
if (DatabaseDescriptor.getNumTokens() == 1)
logger.warn("Generated random token {}. Random tokens will result in an unbalanced ring; see http://wiki.apache.org/cassandra/Operations", bootstrapTokens);
else
logger.info("Generated random tokens. tokens are {}", bootstrapTokens);
}
else
{
bootstrapTokens = new ArrayList<>(initialTokens.size());
for (String token : initialTokens)
bootstrapTokens.add(getTokenFactory().fromString(token));
logger.info("Saved tokens not found. Using configuration value: {}", bootstrapTokens);
}
} }
else else
{ {

View File

@ -232,7 +232,7 @@ public class BootStrapperTest
private void allocateTokensForNode(int vn, String ks, TokenMetadata tm, InetAddress addr) private void allocateTokensForNode(int vn, String ks, TokenMetadata tm, InetAddress addr)
{ {
SummaryStatistics os = TokenAllocation.replicatedOwnershipStats(tm.cloneOnlyTokenMap(), Keyspace.open(ks).getReplicationStrategy(), addr); SummaryStatistics os = TokenAllocation.replicatedOwnershipStats(tm.cloneOnlyTokenMap(), Keyspace.open(ks).getReplicationStrategy(), addr);
Collection<Token> tokens = BootStrapper.allocateTokens(tm, addr, ks, vn); Collection<Token> tokens = BootStrapper.allocateTokens(tm, addr, ks, vn, 0);
assertEquals(vn, tokens.size()); assertEquals(vn, tokens.size());
tm.updateNormalTokens(tokens, addr); tm.updateNormalTokens(tokens, addr);
SummaryStatistics ns = TokenAllocation.replicatedOwnershipStats(tm.cloneOnlyTokenMap(), Keyspace.open(ks).getReplicationStrategy(), addr); SummaryStatistics ns = TokenAllocation.replicatedOwnershipStats(tm.cloneOnlyTokenMap(), Keyspace.open(ks).getReplicationStrategy(), addr);