diff --git a/CHANGES.txt b/CHANGES.txt index 0264dc7457..01f4fd9f4b 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -15,6 +15,7 @@ 2.1.1 + * Add duration mode to cassandra-stress (CASSANDRA-7468) * Add listen_interface and rpc_interface options (CASSANDRA-7417) * Fail to start if commit log replay detects a problem (CASSANDRA-7125) * Improve schema merge performance (CASSANDRA-7444) diff --git a/tools/stress/src/org/apache/cassandra/stress/StressAction.java b/tools/stress/src/org/apache/cassandra/stress/StressAction.java index 2105a7295e..d2253f8f3a 100644 --- a/tools/stress/src/org/apache/cassandra/stress/StressAction.java +++ b/tools/stress/src/org/apache/cassandra/stress/StressAction.java @@ -68,7 +68,8 @@ public class StressAction implements Runnable if (settings.rate.auto) success = runAuto(); else - success = null != run(settings.command.getFactory(settings), settings.rate.threadCount, settings.command.count, output); + success = null != run(settings.command.getFactory(settings), settings.rate.threadCount, settings.command.count, + settings.command.duration, settings.command.durationUnits, output); if (success) output.println("END"); @@ -89,7 +90,7 @@ public class StressAction implements Runnable // we need to warm up all the nodes in the cluster ideally, but we may not be the only stress instance; // so warm up all the nodes we're speaking to only. output.println(String.format("Warming up %s with %d iterations...", single.desc(), iterations)); - run(single, 20, iterations, warmupOutput); + run(single, 20, iterations, 0, null, warmupOutput); } } @@ -105,7 +106,8 @@ public class StressAction implements Runnable { output.println(String.format("Running with %d threadCount", threadCount)); - StressMetrics result = run(settings.command.getFactory(settings), threadCount, settings.command.count, output); + StressMetrics result = run(settings.command.getFactory(settings), threadCount, settings.command.count, + settings.command.duration, settings.command.durationUnits, output); if (result == null) return false; results.add(result); @@ -162,13 +164,14 @@ public class StressAction implements Runnable return improvement / count; } - private StressMetrics run(OpDistributionFactory operations, int threadCount, long opCount, PrintStream output) + private StressMetrics run(OpDistributionFactory operations, int threadCount, long opCount, long duration, TimeUnit durationUnits, PrintStream output) { - output.println(String.format("Running %s with %d threads %s", operations.desc(), threadCount, - opCount > 0 ? " for " + opCount + " iterations" : "until stderr of mean < " + settings.command.targetUncertainty)); + durationUnits != null ? duration + " " + durationUnits.toString().toLowerCase() + : opCount > 0 ? "for " + opCount + " iteration" + : "until stderr of mean < " + settings.command.targetUncertainty)); final WorkQueue workQueue; if (opCount < 0) workQueue = new ContinuousWorkQueue(50); @@ -193,7 +196,12 @@ public class StressAction implements Runnable metrics.start(); - if (opCount <= 0) + if (durationUnits != null) + { + Uninterruptibles.sleepUninterruptibly(duration, durationUnits); + workQueue.stop(); + } + else if (opCount <= 0) { try { diff --git a/tools/stress/src/org/apache/cassandra/stress/StressMetrics.java b/tools/stress/src/org/apache/cassandra/stress/StressMetrics.java index 7e5c1b6310..717dbb866f 100644 --- a/tools/stress/src/org/apache/cassandra/stress/StressMetrics.java +++ b/tools/stress/src/org/apache/cassandra/stress/StressMetrics.java @@ -25,9 +25,10 @@ import java.io.PrintStream; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; +import com.google.common.util.concurrent.Uninterruptibles; import org.apache.commons.lang3.time.DurationFormatUtils; - import org.apache.cassandra.concurrent.NamedThreadFactory; import org.apache.cassandra.stress.util.Timing; import org.apache.cassandra.stress.util.TimingInterval; diff --git a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommand.java b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommand.java index 032f00c3b0..7715e489e6 100644 --- a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommand.java +++ b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommand.java @@ -25,6 +25,7 @@ import java.io.Serializable; import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import org.apache.cassandra.stress.operations.OpDistributionFactory; import org.apache.cassandra.thrift.ConsistencyLevel; @@ -35,6 +36,8 @@ public abstract class SettingsCommand implements Serializable public final Command type; public final long count; + public final long duration; + public final TimeUnit durationUnits; public final int tries; public final boolean ignoreErrors; public final boolean noWarmup; @@ -49,11 +52,12 @@ public abstract class SettingsCommand implements Serializable { this(type, (Options) options, options instanceof Count ? (Count) options : null, + options instanceof Duration ? (Duration) options : null, options instanceof Uncertainty ? (Uncertainty) options : null ); } - public SettingsCommand(Command type, Options options, Count count, Uncertainty uncertainty) + public SettingsCommand(Command type, Options options, Count count, Duration duration, Uncertainty uncertainty) { this.type = type; this.tries = Math.max(1, Integer.parseInt(options.retries.value()) + 1); @@ -63,6 +67,30 @@ public abstract class SettingsCommand implements Serializable if (count != null) { this.count = Long.parseLong(count.count.value()); + this.duration = 0; + this.durationUnits = null; + this.targetUncertainty = -1; + this.minimumUncertaintyMeasurements = -1; + this.maximumUncertaintyMeasurements = -1; + } + else if (duration != null) + { + this.count = -1; + this.duration = Long.parseLong(duration.duration.value().substring(0, duration.duration.value().length() - 1)); + switch (duration.duration.value().toLowerCase().charAt(duration.duration.value().length() - 1)) + { + case 's': + this.durationUnits = TimeUnit.SECONDS; + break; + case 'm': + this.durationUnits = TimeUnit.MINUTES; + break; + case 'h': + this.durationUnits = TimeUnit.HOURS; + break; + default: + throw new IllegalStateException(); + } this.targetUncertainty = -1; this.minimumUncertaintyMeasurements = -1; this.maximumUncertaintyMeasurements = -1; @@ -70,6 +98,8 @@ public abstract class SettingsCommand implements Serializable else { this.count = -1; + this.duration = 0; + this.durationUnits = null; this.targetUncertainty = Double.parseDouble(uncertainty.uncertainty.value()); this.minimumUncertaintyMeasurements = Integer.parseInt(uncertainty.minMeasurements.value()); this.maximumUncertaintyMeasurements = Integer.parseInt(uncertainty.maxMeasurements.value()); @@ -97,6 +127,16 @@ public abstract class SettingsCommand implements Serializable } } + static class Duration extends Options + { + final OptionSimple duration = new OptionSimple("duration=", "[0-9]+[smh]", null, "Time to run in (in seconds, minutes or hours)", true); + @Override + public List options() + { + return Arrays.asList(duration, retries, ignoreErrors, consistencyLevel, atOnce); + } + } + static class Uncertainty extends Options { final OptionSimple uncertainty = new OptionSimple("err<", "0\\.[0-9]+", "0.02", "Run until the standard error of the mean is below this fraction", false); @@ -138,5 +178,37 @@ public abstract class SettingsCommand implements Serializable return null; } -} +/* static SettingsCommand build(Command type, String[] params) + { + GroupedOptions options = GroupedOptions.select(params, new Count(), new Duration(), new Uncertainty()); + if (options == null) + { + printHelp(type); + System.out.println("Invalid " + type + " options provided, see output for valid options"); + System.exit(1); + } + return new SettingsCommand(type, options); + }*/ + static void printHelp(Command type) + { + printHelp(type.toString().toLowerCase()); + } + + static void printHelp(String type) + { + GroupedOptions.printOptions(System.out, type.toLowerCase(), new Uncertainty(), new Count(), new Duration()); + } + + static Runnable helpPrinter(final Command type) + { + return new Runnable() + { + @Override + public void run() + { + printHelp(type); + } + }; + } +} diff --git a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefined.java b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefined.java index ac113d17b7..64dc054671 100644 --- a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefined.java +++ b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefined.java @@ -35,6 +35,7 @@ import org.apache.cassandra.stress.operations.FixedOpDistribution; import org.apache.cassandra.stress.operations.OpDistribution; import org.apache.cassandra.stress.operations.OpDistributionFactory; import org.apache.cassandra.stress.operations.predefined.PredefinedOperation; +import org.apache.cassandra.stress.settings.SettingsCommandPreDefinedMixed.Options; import org.apache.cassandra.stress.util.Timer; // Settings unique to the mixed command type @@ -111,7 +112,8 @@ public class SettingsCommandPreDefined extends SettingsCommand { GroupedOptions options = GroupedOptions.select(params, new Options(new Uncertainty()), - new Options(new Count())); + new Options(new Count()), + new Options(new Duration())); if (options == null) { printHelp(type); @@ -128,7 +130,7 @@ public class SettingsCommandPreDefined extends SettingsCommand static void printHelp(String type) { - GroupedOptions.printOptions(System.out, type.toLowerCase(), new Uncertainty(), new Count()); + GroupedOptions.printOptions(System.out, type.toLowerCase(), new Uncertainty(), new Count(), new Duration()); } static Runnable helpPrinter(final Command type) diff --git a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefinedMixed.java b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefinedMixed.java index e5d4f8095b..858534beb2 100644 --- a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefinedMixed.java +++ b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandPreDefinedMixed.java @@ -120,7 +120,8 @@ public class SettingsCommandPreDefinedMixed extends SettingsCommandPreDefined { GroupedOptions options = GroupedOptions.select(params, new Options(new SettingsCommand.Uncertainty()), - new Options(new SettingsCommand.Count())); + new Options(new SettingsCommand.Count()), + new Options(new SettingsCommand.Duration())); if (options == null) { printHelp(); @@ -134,7 +135,8 @@ public class SettingsCommandPreDefinedMixed extends SettingsCommandPreDefined { GroupedOptions.printOptions(System.out, "mixed", new Options(new SettingsCommand.Uncertainty()), - new Options(new SettingsCommand.Count())); + new Options(new SettingsCommand.Count()), + new Options(new SettingsCommand.Duration())); } public static Runnable helpPrinter() diff --git a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandUser.java b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandUser.java index f36296e3c5..6463783547 100644 --- a/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandUser.java +++ b/tools/stress/src/org/apache/cassandra/stress/settings/SettingsCommandUser.java @@ -118,7 +118,8 @@ public class SettingsCommandUser extends SettingsCommand { GroupedOptions.printOptions(System.out, "user", new Options(new Uncertainty()), - new Options(new Count())); + new Options(new Count()), + new Options(new Duration())); } public static Runnable helpPrinter()