diff --git a/CHANGES.txt b/CHANGES.txt index b4be8d01b1..3056988939 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -17,7 +17,13 @@ * fixes for cache save/load (CASSANDRA-2172, -2174) * Handle whole-row deletions in CFOutputFormat (CASSANDRA-2014) * Make memtable_flush_writers flush in parallel (CASSANDRA-2178) + * make key cache preheating default to false; enable with + -Dcompaction_preheat_key_cache=true (CASSANDRA-2175) + * refactor stress.py to have only one copy of the format string + used for creating row keys (CASSANDRA-2108) + * validate index names for \w+ (CASSANDRA-2196) * Fix Cassandra cli to respect timeout if schema does not settle (CASSANDRA-2187) + * update memtable_throughput to be a long (CASSANDRA-2158) 0.7.2 @@ -92,6 +98,7 @@ * bound hints CF throughput between 32M and 256M (CASSANDRA-2148) * continue starting when invalid saved cache entries are encountered (CASSANDRA-2076) + * add max_hint_window_in_ms option (CASSANDRA-1459) 0.7.0-final diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index f4eaad11f8..e0036492a4 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -231,6 +231,11 @@ column_index_size_in_kb: 64 # will be logged specifying the row key. in_memory_compaction_limit_in_mb: 64 +# Track cached row keys during compaction, and re-cache their new +# positions in the compacted sstable. Disable if you use really large +# key caches. +compaction_preheat_key_cache: true + # Time to wait for a reply from other nodes before failing the command rpc_timeout_in_ms: 10000 diff --git a/contrib/py_stress/stress.py b/contrib/py_stress/stress.py index 251214e782..7d89fbff4f 100644 --- a/contrib/py_stress/stress.py +++ b/contrib/py_stress/stress.py @@ -127,6 +127,9 @@ if options.nodefile != None: with open(options.nodefile) as f: nodes = [n.strip() for n in f.readlines() if len(n.strip()) > 0] +#format string for keys +fmt = '%0' + str(len(str(total_keys))) + 'd' + # a generator that generates all keys according to a bell curve centered # around the middle of the keys generated (0..total_keys). Remember that # about 68% of keys will be within stdev away from the mean and @@ -148,7 +151,6 @@ def generate_values(): return values def key_generator_gauss(): - fmt = '%0' + str(len(str(total_keys))) + 'd' while True: guess = gauss(mean, stdev) if 0 <= guess < total_keys: @@ -157,7 +159,6 @@ def key_generator_gauss(): # a generator that will generate all keys w/ equal probability. this is the # worst case for caching. def key_generator_random(): - fmt = '%0' + str(len(str(total_keys))) + 'd' return fmt % randint(0, total_keys - 1) key_generator = key_generator_gauss @@ -220,7 +221,6 @@ class Inserter(Operation): def run(self): values = generate_values() columns = [Column('C' + str(j), 'unset', time.time() * 1000000) for j in xrange(columns_per_key)] - fmt = '%0' + str(len(str(total_keys))) + 'd' if 'super' == options.cftype: supers = [SuperColumn('S' + str(j), columns) for j in xrange(supers_per_key)] for i in self.range: @@ -295,7 +295,6 @@ class RangeSlicer(Operation): end = self.range[-1] current = begin last = current + options.rangecount - fmt = '%0' + str(len(str(total_keys))) + 'd' p = SlicePredicate(slice_range=SliceRange('', '', False, columns_per_key)) if 'super' == options.cftype: while current < end: @@ -348,7 +347,6 @@ class RangeSlicer(Operation): # from the thread's appointed range class IndexedRangeSlicer(Operation): def run(self): - fmt = '%0' + str(len(str(total_keys))) + 'd' p = SlicePredicate(slice_range=SliceRange('', '', False, columns_per_key)) values = generate_values() parent = ColumnParent('Standard1') diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java index 73f7413d5c..7138c6be9a 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java @@ -19,9 +19,7 @@ package org.apache.cassandra.contrib.stress; import java.io.*; import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.List; +import java.util.*; import java.util.concurrent.atomic.AtomicIntegerArray; import java.util.concurrent.atomic.AtomicLongArray; @@ -29,6 +27,7 @@ import org.apache.commons.cli.*; import org.apache.cassandra.db.ColumnFamilyType; import org.apache.cassandra.thrift.*; +import org.apache.commons.lang.StringUtils; import org.apache.thrift.protocol.TBinaryProtocol; import org.apache.thrift.transport.TFramedTransport; import org.apache.thrift.transport.TSocket; @@ -62,12 +61,15 @@ public class Session availableOptions.addOption("o", "operation", true, "Operation to perform (INSERT, READ, RANGE_SLICE, INDEXED_RANGE_SLICE, MULTI_GET), default:INSERT"); availableOptions.addOption("u", "supercolumns", true, "Number of super columns per key, default:1"); availableOptions.addOption("y", "family-type", true, "Column Family Type (Super, Standard), default:Standard"); - availableOptions.addOption("k", "keep-going", false, "Ignore errors inserting or reading, default:false"); + availableOptions.addOption("K", "keep-trying", true, "Retry on-going operation N times (in case of failure). positive integer, default:10"); + availableOptions.addOption("k", "keep-going", false, "Ignore errors inserting or reading (when set, --keep-trying has no effect), default:false"); availableOptions.addOption("i", "progress-interval", true, "Progress Report Interval (seconds), default:10"); availableOptions.addOption("g", "keys-per-call", true, "Number of keys to get_range_slices or multiget per call, default:1000"); availableOptions.addOption("l", "replication-factor", true, "Replication Factor to use when creating needed column families, default:1"); availableOptions.addOption("e", "consistency-level", true, "Consistency Level to use (ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL, ANY), default:ONE"); availableOptions.addOption("x", "create-index", true, "Type of index to create on needed column families (KEYS)"); + availableOptions.addOption("R", "replication-strategy", true, "Replication strategy to use (only on insert if keyspace does not exist), default:org.apache.cassandra.locator.SimpleStrategy"); + availableOptions.addOption("O", "strategy-properties", true, "Replication strategy properties in the following format :,:,..."); } private int numKeys = 1000 * 1000; @@ -79,13 +81,14 @@ public class Session private String[] nodes = new String[] { "127.0.0.1" }; private boolean random = false; private boolean unframed = false; - private boolean ignoreErrors = false; + private int retryTimes = 10; private int port = 9160; private int superColumns = 1; private int progressInterval = 10; private int keysPerCall = 1000; private int replicationFactor = 1; + private boolean ignoreErrors = false; private PrintStream out = System.out; @@ -93,6 +96,9 @@ public class Session private Stress.Operation operation = Stress.Operation.INSERT; private ColumnFamilyType columnFamilyType = ColumnFamilyType.Standard; private ConsistencyLevel consistencyLevel = ConsistencyLevel.ONE; + private String replicationStrategy = "org.apache.cassandra.locator.SimpleStrategy"; + private Map replicationStrategyOptions = new HashMap(); + // required by Gaussian distribution. protected int mean; @@ -185,8 +191,21 @@ public class Session if (cmd.hasOption("y")) columnFamilyType = ColumnFamilyType.valueOf(cmd.getOptionValue("y")); + if (cmd.hasOption("K")) + { + retryTimes = Integer.valueOf(cmd.getOptionValue("K")); + + if (retryTimes <= 0) + { + throw new RuntimeException("--keep-trying option value should be > 0"); + } + } + if (cmd.hasOption("k")) + { + retryTimes = 1; ignoreErrors = true; + } if (cmd.hasOption("i")) progressInterval = Integer.parseInt(cmd.getOptionValue("i")); @@ -202,6 +221,24 @@ public class Session if (cmd.hasOption("x")) indexType = IndexType.valueOf(cmd.getOptionValue("x").toUpperCase()); + + if (cmd.hasOption("R")) + replicationStrategy = cmd.getOptionValue("R"); + + if (cmd.hasOption("O")) + { + String[] pairs = StringUtils.split(cmd.getOptionValue("O"), ','); + + for (String pair : pairs) + { + String[] keyAndValue = StringUtils.split(pair, ':'); + + if (keyAndValue.length != 2) + throw new RuntimeException("Invalid --strategy-properties value."); + + replicationStrategyOptions.put(keyAndValue[0], keyAndValue[1]); + } + } } catch (ParseException e) { @@ -276,6 +313,11 @@ public class Session return consistencyLevel; } + public int getRetryTimes() + { + return retryTimes; + } + public boolean ignoreErrors() { return ignoreErrors; @@ -337,8 +379,14 @@ public class Session CfDef superCfDef = new CfDef("Keyspace1", "Super1").setColumn_metadata(Arrays.asList(superSubColumn)).setColumn_type("Super"); keyspace.setName("Keyspace1"); - keyspace.setStrategy_class("org.apache.cassandra.locator.SimpleStrategy"); + keyspace.setStrategy_class(replicationStrategy); keyspace.setReplication_factor(replicationFactor); + + if (!replicationStrategyOptions.isEmpty()) + { + keyspace.setStrategy_options(replicationStrategyOptions); + } + keyspace.setCf_defs(new ArrayList(Arrays.asList(standardCfDef, superCfDef))); Cassandra.Client client = getClient(false); diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java index 6891921fe9..451ca9871b 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java @@ -63,21 +63,33 @@ public class IndexedRangeSlicer extends OperationThread List results = null; long start = System.currentTimeMillis(); - try + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { - results = client.get_indexed_slices(parent, clause, predicate, session.getConsistencyLevel()); + if (success) + break; - if (results.size() == 0) + try { - System.err.printf("No indexed values from offset received: %s%n", startOffset); - - if (!session.ignoreErrors()) - break; + results = client.get_indexed_slices(parent, clause, predicate, session.getConsistencyLevel()); + success = (results.size() != 0); + } + catch (Exception e) + { + exceptionMessage = getExceptionMessage(e); + success = false; } } - catch (Exception e) + + if (!success) { - System.err.printf("Error on get_indexed_slices call for offset %s - %s%n", startOffset, getExceptionMessage(e)); + System.err.printf("Thread [%d] retried %d times - error on calling get_indexed_slices for offset %s %s%n", + index, + session.getRetryTimes(), + startOffset, + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) return; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java index 95c26facc5..279395f77e 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java @@ -64,7 +64,8 @@ public class Inserter extends OperationThread for (int i : range) { - ByteBuffer key = ByteBuffer.wrap(String.format(format, i).getBytes()); + String rawKey = String.format(format, i); + ByteBuffer key = ByteBuffer.wrap(rawKey.getBytes()); Map>> record = new HashMap>>(); record.put(key, session.getColumnFamilyType() == ColumnFamilyType.Super @@ -78,23 +79,35 @@ public class Inserter extends OperationThread long start = System.currentTimeMillis(); - try - { - client.batch_mutate(record, session.getConsistencyLevel()); - } - catch (Exception e) + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { + if (success) + break; + try { - System.err.printf("Error while inserting key %s - %s%n", ByteBufferUtil.string(key), getExceptionMessage(e)); + client.batch_mutate(record, session.getConsistencyLevel()); + success = true; } - catch (CharacterCodingException e1) + catch (Exception e) { - throw new AssertionError(e1); // keys are valid strings + exceptionMessage = getExceptionMessage(e); + success = false; } + } + + if (!success) + { + System.err.printf("Thread [%d] retried %d times - error inserting key %s %s%n", index, + session.getRetryTimes(), + rawKey, + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) - return; + break; } session.operationCount.getAndIncrement(index); diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java index f0bbe76666..d842daf345 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java @@ -55,21 +55,32 @@ public class MultiGetter extends OperationThread long start = System.currentTimeMillis(); - try + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { - results = client.multiget_slice(keys, parent, predicate, session.getConsistencyLevel()); + if (success) + break; - if (results.size() == 0) + try { - System.err.printf("Keys %s were not found.%n", keys); - - if (!session.ignoreErrors()) - break; + results = client.multiget_slice(keys, parent, predicate, session.getConsistencyLevel()); + success = (results.size() != 0); + } + catch (Exception e) + { + exceptionMessage = getExceptionMessage(e); } } - catch (Exception e) + + if (!success) { - System.err.printf("Error on multiget_slice call - %s%n", getExceptionMessage(e)); + System.err.printf("Thread [%d] retried %d times - error on calling multiget_slice for keys %s %s%n", + index, + session.getRetryTimes(), + keys, + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) return; @@ -93,21 +104,33 @@ public class MultiGetter extends OperationThread long start = System.currentTimeMillis(); - try + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { - results = client.multiget_slice(keys, parent, predicate, session.getConsistencyLevel()); + if (success) + break; - if (results.size() == 0) + try { - System.err.printf("Keys %s were not found.%n", keys); - - if (!session.ignoreErrors()) - break; + results = client.multiget_slice(keys, parent, predicate, session.getConsistencyLevel()); + success = (results.size() != 0); + } + catch (Exception e) + { + exceptionMessage = getExceptionMessage(e); + success = false; } } - catch (Exception e) + + if (!success) { - System.err.printf("Error on multiget_slice call - %s%n", getExceptionMessage(e)); + System.err.printf("Thread [%d] retried %d times - error on calling multiget_slice for keys %s %s%n", + index, + session.getRetryTimes(), + keys, + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) return; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java index 8cdf6d6f8b..644410bd0c 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java @@ -64,21 +64,31 @@ public class RangeSlicer extends OperationThread long startTime = System.currentTimeMillis(); - try + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { - slices = client.get_range_slices(parent, predicate, range, session.getConsistencyLevel()); - - if (slices.size() == 0) + try { - System.err.printf("Range %s->%s not found in Super Column %s.%n", new String(start), new String(end), superColumnName); - - if (!session.ignoreErrors()) - break; + slices = client.get_range_slices(parent, predicate, range, session.getConsistencyLevel()); + success = (slices.size() != 0); + } + catch (Exception e) + { + exceptionMessage = getExceptionMessage(e); + success = false; } } - catch (Exception e) + + if (!success) { - System.err.printf("Error while reading Super Column %s - %s%n", superColumnName, getExceptionMessage(e)); + System.err.printf("Thread [%d] retried %d times - error on calling get_range_slices for range %s->%s %s%n", + index, + session.getRetryTimes(), + new String(start), + new String(end), + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) return; @@ -107,21 +117,34 @@ public class RangeSlicer extends OperationThread long startTime = System.currentTimeMillis(); - try + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { - slices = client.get_range_slices(parent, predicate, range, session.getConsistencyLevel()); + if (success) + break; - if (slices.size() == 0) + try { - System.err.printf("Range %s->%s not found.%n", String.format(format, current), String.format(format, last)); - - if (!session.ignoreErrors()) - break; + slices = client.get_range_slices(parent, predicate, range, session.getConsistencyLevel()); + success = (slices.size() != 0); + } + catch (Exception e) + { + exceptionMessage = getExceptionMessage(e); + success = false; } } - catch (Exception e) + + if (!success) { - System.err.printf("Error while reading range %s->%s - %s%n", String.format(format, current), String.format(format, last), getExceptionMessage(e)); + System.err.printf("Thread [%d] retried %d times - error on calling get_indexed_slices for range %s->%s %s%n", + index, + session.getRetryTimes(), + new String(start), + new String(end), + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) return; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java index 5eb6f0c2d3..acbbdcc6dc 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java @@ -22,6 +22,7 @@ import org.apache.cassandra.db.ColumnFamilyType; import org.apache.cassandra.thrift.*; import org.apache.cassandra.utils.ByteBufferUtil; +import java.io.IOException; import java.lang.AssertionError; import java.nio.ByteBuffer; import java.nio.charset.CharacterCodingException; @@ -61,7 +62,8 @@ public class Reader extends OperationThread { for (int i = 0; i < session.getKeysPerThread(); i++) { - ByteBuffer key = ByteBuffer.wrap(generateKey()); + byte[] rawKey = generateKey(); + ByteBuffer key = ByteBuffer.wrap(rawKey); for (int j = 0; j < session.getSuperColumns(); j++) { @@ -70,32 +72,36 @@ public class Reader extends OperationThread long start = System.currentTimeMillis(); - try - { - List columns; - columns = client.get_slice(key, parent, predicate, session.getConsistencyLevel()); + boolean success = false; + String exceptionMessage = null; - if (columns.size() == 0) - { - System.err.printf("Key %s not found in Super Column %s.%n", ByteBufferUtil.string(key), superColumn); - - if (!session.ignoreErrors()) - break; - } - } - catch (Exception e) + for (int t = 0; t < session.getRetryTimes(); t++) { + if (success) + break; + try { - System.err.printf("Error while reading Super Column %s key %s - %s%n", superColumn, ByteBufferUtil.string(key), getExceptionMessage(e)); + List columns; + columns = client.get_slice(key, parent, predicate, session.getConsistencyLevel()); + success = (columns.size() != 0); } - catch (CharacterCodingException e1) + catch (Exception e) { - throw new AssertionError(e1); // keys are valid string + exceptionMessage = getExceptionMessage(e); + success = false; } + } + + if (!success) + { + System.err.printf("Thread [%d] retried %d times - error reading key %s %s%n", index, + session.getRetryTimes(), + new String(rawKey), + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) - break; + return; } session.operationCount.getAndIncrement(index); @@ -116,25 +122,36 @@ public class Reader extends OperationThread long start = System.currentTimeMillis(); - try + boolean success = false; + String exceptionMessage = null; + + for (int t = 0; t < session.getRetryTimes(); t++) { - List columns; - columns = client.get_slice(keyBuffer, parent, predicate, session.getConsistencyLevel()); + if (success) + break; - if (columns.size() == 0) + try { - System.err.println(String.format("Key %s not found.", new String(key))); - - if (!session.ignoreErrors()) - break; + List columns; + columns = client.get_slice(keyBuffer, parent, predicate, session.getConsistencyLevel()); + success = (columns.size() != 0); + } + catch (Exception e) + { + exceptionMessage = getExceptionMessage(e); + success = false; } } - catch (Exception e) + + if (!success) { - System.err.printf("Error while reading key %s - %s%n", new String(key), getExceptionMessage(e)); + System.err.printf("Thread [%d] retried %d times - error reading key %s %s%n", index, + session.getRetryTimes(), + new String(key), + (exceptionMessage == null) ? "" : "(" + exceptionMessage + ")"); if (!session.ignoreErrors()) - break; + return; } session.operationCount.getAndIncrement(index); diff --git a/src/java/org/apache/cassandra/cli/CliMain.java b/src/java/org/apache/cassandra/cli/CliMain.java index 3b24c70f72..4331895e47 100644 --- a/src/java/org/apache/cassandra/cli/CliMain.java +++ b/src/java/org/apache/cassandra/cli/CliMain.java @@ -90,32 +90,6 @@ public class CliMain thriftClient = cassandraClient; cliClient = new CliClient(sessionState, thriftClient); - if (sessionState.keyspace != null) - { - try - { - sessionState.keyspace = CliCompiler.getKeySpace(sessionState.keyspace, thriftClient.describe_keyspaces());; - thriftClient.set_keyspace(sessionState.keyspace); - cliClient.setKeySpace(sessionState.keyspace); - updateCompletor(CliUtils.getCfNamesByKeySpace(cliClient.getKSMetaData(sessionState.keyspace))); - } - catch (InvalidRequestException e) - { - sessionState.err.println("Keyspace " + sessionState.keyspace + " not found"); - return; - } - catch (TException e) - { - sessionState.err.println("Did you specify 'keyspace'?"); - return; - } - catch (NotFoundException e) - { - sessionState.err.println("Keyspace " + sessionState.keyspace + " not found"); - return; - } - } - if ((sessionState.username != null) && (sessionState.password != null)) { // Authenticate @@ -149,6 +123,32 @@ public class CliMain } } + if (sessionState.keyspace != null) + { + try + { + sessionState.keyspace = CliCompiler.getKeySpace(sessionState.keyspace, thriftClient.describe_keyspaces());; + thriftClient.set_keyspace(sessionState.keyspace); + cliClient.setKeySpace(sessionState.keyspace); + updateCompletor(CliUtils.getCfNamesByKeySpace(cliClient.getKSMetaData(sessionState.keyspace))); + } + catch (InvalidRequestException e) + { + sessionState.err.println("Keyspace " + sessionState.keyspace + " not found"); + return; + } + catch (TException e) + { + sessionState.err.println("Did you specify 'keyspace'?"); + return; + } + catch (NotFoundException e) + { + sessionState.err.println("Keyspace " + sessionState.keyspace + " not found"); + return; + } + } + // Lookup the cluster name, this is to make it clear which cluster the user is connected to String clusterName; diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index ede3cd3690..915bb58f5e 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -110,6 +110,7 @@ public class Config public Double reduce_cache_sizes_at = 1.0; public double reduce_cache_capacity_to = 0.6; public int hinted_handoff_throttle_delay_in_ms = 0; + public boolean compaction_preheat_key_cache = true; public static enum CommitLogSync { periodic, diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 4f1111cbc4..9b3f4d3e8e 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -1196,4 +1196,9 @@ public class DatabaseDescriptor { return conf.hinted_handoff_throttle_delay_in_ms; } + + public static boolean getPreheatKeyCache() + { + return conf.compaction_preheat_key_cache; + } } diff --git a/src/java/org/apache/cassandra/db/CompactionManager.java b/src/java/org/apache/cassandra/db/CompactionManager.java index 6fa7b8e020..54bf6f2b7f 100644 --- a/src/java/org/apache/cassandra/db/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/CompactionManager.java @@ -431,12 +431,15 @@ public class CompactionManager implements CompactionManagerMBean long position = writer.append(row); totalkeysWritten++; - for (SSTableReader sstable : sstables) + if (DatabaseDescriptor.getPreheatKeyCache()) { - if (sstable.getCachedPosition(row.key) != null) + for (SSTableReader sstable : sstables) { - cachedKeys.put(row.key, position); - break; + if (sstable.getCachedPosition(row.key) != null) + { + cachedKeys.put(row.key, position); + break; + } } } } @@ -448,7 +451,7 @@ public class CompactionManager implements CompactionManagerMBean SSTableReader ssTable = writer.closeAndOpenReader(getMaxDataAge(sstables)); cfs.replaceCompactedSSTables(sstables, Arrays.asList(ssTable)); - for (Entry entry : cachedKeys.entrySet()) + for (Entry entry : cachedKeys.entrySet()) // empty if preheat is off ssTable.cacheKey(entry.getKey(), entry.getValue()); submitMinorIfNeeded(cfs); diff --git a/test/unit/org/apache/cassandra/streaming/StreamingTransferTest.java b/test/unit/org/apache/cassandra/streaming/StreamingTransferTest.java index 46bc59c8b0..abf867b771 100644 --- a/test/unit/org/apache/cassandra/streaming/StreamingTransferTest.java +++ b/test/unit/org/apache/cassandra/streaming/StreamingTransferTest.java @@ -28,6 +28,7 @@ import java.util.*; import org.apache.cassandra.CleanupHelper; import org.apache.cassandra.Util; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.db.*; import org.apache.cassandra.db.columniterator.IdentityQueryFilter; import org.apache.cassandra.db.filter.IFilter; @@ -155,4 +156,49 @@ public class StreamingTransferTest extends CleanupHelper assert null != cfstore.getColumnFamily(QueryFilter.getIdentityFilter(Util.dk("test"), new QueryPath("Standard1"))); assert null != cfstore.getColumnFamily(QueryFilter.getIdentityFilter(Util.dk("transfer1"), new QueryPath("Standard1"))); } + + @Test + public void testTransferOfMultipleColumnFamilies() throws Exception + { + String keyspace = "Keyspace1"; + IPartitioner p = StorageService.getPartitioner(); + String[] columnFamilies = new String[] { "Standard1", "Standard2", "Standard3" }; + List ssTableReaders = new ArrayList(); + + // ranges to transfer + List ranges = new ArrayList(); + + for (String cf : columnFamilies) + { + Set content = new HashSet(); + + content.add("data-" + cf + "-1"); + content.add("data-" + cf + "-2"); + content.add("data-" + cf + "-3"); + + SSTableUtils.Context context = SSTableUtils.prepare().ks(keyspace).cf(cf); + + ssTableReaders.add(context.write(content)); + ranges.add(new Range(p.getMinimumToken(), p.getToken(ByteBufferUtil.bytes("data-" + cf + "-3")))); + } + + StreamOutSession session = StreamOutSession.create(keyspace, LOCAL, null); + StreamOut.transferSSTables(session, ssTableReaders, ranges); + + session.await(); + + for (String cf : columnFamilies) + { + ColumnFamilyStore store = Table.open(keyspace).getColumnFamilyStore(cf); + List rows = Util.getRangeSlice(store); + + assert rows.size() >= 3; + + for (int i = 0; i < 3; i++) + { + String expectedKey = "data-" + cf + "-" + (i + 1); + assertEquals(p.decorateKey(ByteBufferUtil.bytes(expectedKey)), rows.get(i).key); + } + } + } }