diff --git a/src/java/org/apache/cassandra/locator/TokenMetadata.java b/src/java/org/apache/cassandra/locator/TokenMetadata.java index 164d80fde6..771ce2f727 100644 --- a/src/java/org/apache/cassandra/locator/TokenMetadata.java +++ b/src/java/org/apache/cassandra/locator/TokenMetadata.java @@ -287,4 +287,9 @@ public class TokenMetadata { return getEndPoint(getSuccessor(getToken(endPoint))); } + + public void clearPendingRanges() + { + pendingRanges.clear(); + } } diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index bc23170682..d4beeb8c9a 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -1046,4 +1046,9 @@ public final class StorageService implements IEndPointStateChangeSubscriber, Sto { return replicationStrategy_; } + + public void cancelPendingRanges() + { + tokenMetadata_.clearPendingRanges(); + } } diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java index f4a991319e..3cb2001f50 100644 --- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java +++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java @@ -131,6 +131,13 @@ public interface StorageServiceMBean */ public void loadBalance() throws IOException, InterruptedException; + /** + * cancel writes to nodes that are set to be changing ranges. + * Only do this if the reason for the range changes no longer exists + * (e.g., a bootstrapping node was killed or crashed.) + */ + public void cancelPendingRanges(); + /** set the logging level at runtime */ public void setLog4jLevel(String classQualifier, String level); } diff --git a/src/java/org/apache/cassandra/tools/NodeProbe.java b/src/java/org/apache/cassandra/tools/NodeProbe.java index 96d80bf8df..347e2c16c2 100644 --- a/src/java/org/apache/cassandra/tools/NodeProbe.java +++ b/src/java/org/apache/cassandra/tools/NodeProbe.java @@ -393,6 +393,11 @@ public class NodeProbe ssProxy.move(newToken); } + public void cancelPendingRanges() + { + ssProxy.cancelPendingRanges(); + } + /** * Print out the size of the queues in the thread pools * @@ -488,7 +493,7 @@ public class NodeProbe HelpFormatter hf = new HelpFormatter(); String header = String.format( "%nAvailable commands: ring, info, cleanup, compact, cfstats, snapshot [name], clearsnapshot, " + - "tpstats, flush, decommission, move, loadbalance, " + + "tpstats, flush, decommission, move, loadbalance, cancelpending, " + " getcompactionthreshold, setcompactionthreshold [minthreshold] ([maxthreshold])"); String usage = String.format("java %s -host %n", NodeProbe.class.getName()); hf.printHelp(usage, "", options, header); @@ -563,6 +568,10 @@ public class NodeProbe } probe.move(arguments[1]); } + else if (cmdName.equals("cancelpending")) + { + probe.cancelPendingRanges(); + } else if (cmdName.equals("snapshot")) { String snapshotName = "";