diff --git a/CHANGES.txt b/CHANGES.txt index 884ee1ed7f..da076f81ac 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -40,6 +40,12 @@ * add ant-optional as dependence for the debian package (CASSANDRA-2164) * add option to specify limit for get_slice in the CLI (CASSANDRA-2646) * decrease HH page size (CASSANDRA-2832) + * reset cli keyspace after dropping the current one (CASSANDRA-2763) + * add KeyRange option to Hadoop inputformat (CASSANDRA-1125) + * fix protocol versioning (CASSANDRA-2818, 2860) + * support spaces in path to log4j configuration (CASSANDRA-2383) + * avoid including inferred types in CF update (CASSANDRA-2809) + * fix JMX bulkload call (CASSANDRA-2908) 0.8.1 @@ -232,6 +238,7 @@ * add a server-wide cap on measured memtable memory usage and aggressively flush to keep under that threshold (CASSANDRA-2006) * add unified UUIDType (CASSANDRA-2233) + * add off-heap row cache support (CASSANDRA-1969) 0.7.5 diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 625f79a7ba..1901ed25ee 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -80,13 +80,16 @@ saved_caches_directory: /var/lib/cassandra/saved_caches # commitlog_sync may be either "periodic" or "batch." # When in batch mode, Cassandra won't ack writes until the commit log # has been fsynced to disk. It will wait up to -# CommitLogSyncBatchWindowInMS milliseconds for other writes, before +# commitlog_sync_batch_window_in_ms milliseconds for other writes, before # performing the sync. -commitlog_sync: periodic - +# +# commitlog_sync: batch +# commitlog_sync_batch_window_in_ms: 50 +# # the other option is "periodic" where writes may be acked immediately # and the CommitLog is simply synced every commitlog_sync_period_in_ms # milliseconds. +commitlog_sync: periodic commitlog_sync_period_in_ms: 10000 # any class that implements the SeedProvider interface and has a constructor that takes a Map of diff --git a/examples/client_only/conf/cassandra.yaml b/examples/client_only/conf/cassandra.yaml index 97e3853e4b..2d92794faf 100644 --- a/examples/client_only/conf/cassandra.yaml +++ b/examples/client_only/conf/cassandra.yaml @@ -77,9 +77,6 @@ commitlog_directory: /var/lib/cassandra/commitlog # saved caches saved_caches_directory: /var/lib/cassandra/saved_caches -# Size to allow commitlog to grow to before creating a new segment -commitlog_rotation_threshold_in_mb: 128 - # commitlog_sync may be either "periodic" or "batch." # When in batch mode, Cassandra won't ack writes until the commit log # has been fsynced to disk. It will wait up to diff --git a/src/java/org/apache/cassandra/cli/CliClient.java b/src/java/org/apache/cassandra/cli/CliClient.java index 7ef940a7e9..878bf7462b 100644 --- a/src/java/org/apache/cassandra/cli/CliClient.java +++ b/src/java/org/apache/cassandra/cli/CliClient.java @@ -350,7 +350,6 @@ public class CliClient Tree columnFamilySpec = statement.getChild(0); - String key = CliCompiler.getKey(columnFamilySpec); String columnFamily = CliCompiler.getColumnFamily(columnFamilySpec, keyspacesMap.get(keySpace).cf_defs); int columnSpecCnt = CliCompiler.numColumnSpecifiers(columnFamilySpec); @@ -358,14 +357,19 @@ public class CliClient if (columnSpecCnt != 0) { - byte[] superColumn = columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), columnFamily); + Tree columnTree = columnFamilySpec.getChild(2); + + byte[] superColumn = (columnTree.getType() == CliParser.FUNCTION_CALL) + ? convertValueByFunction(columnTree, null, null).array() + : columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), columnFamily); + colParent = new ColumnParent(columnFamily).setSuper_column(superColumn); } SliceRange range = new SliceRange(ByteBufferUtil.EMPTY_BYTE_BUFFER, ByteBufferUtil.EMPTY_BYTE_BUFFER, false, Integer.MAX_VALUE); SlicePredicate predicate = new SlicePredicate().setColumn_names(null).setSlice_range(range); - int count = thriftClient.get_count(ByteBufferUtil.bytes(key), colParent, predicate, consistencyLevel); + int count = thriftClient.get_count(getKeyAsBytes(columnFamily, columnFamilySpec.getChild(1)), colParent, predicate, consistencyLevel); sessionState.out.printf("%d columns%n", count); } @@ -377,13 +381,14 @@ public class CliClient Tree columnFamilySpec = statement.getChild(0); - String key = CliCompiler.getKey(columnFamilySpec); String columnFamily = CliCompiler.getColumnFamily(columnFamilySpec, keyspacesMap.get(keySpace).cf_defs); + CfDef cfDef = getCfDef(columnFamily); + + ByteBuffer key = getKeyAsBytes(columnFamily, columnFamilySpec.getChild(1)); int columnSpecCnt = CliCompiler.numColumnSpecifiers(columnFamilySpec); byte[] superColumnName = null; byte[] columnName = null; - CfDef cfDef = getCfDef(columnFamily); boolean isSuper = cfDef.column_type.equals("Super"); if ((columnSpecCnt < 0) || (columnSpecCnt > 2)) @@ -391,20 +396,42 @@ public class CliClient sessionState.out.println("Invalid row, super column, or column specification."); return; } - + + Tree columnTree = (columnSpecCnt >= 1) + ? columnFamilySpec.getChild(2) + : null; + + Tree subColumnTree = (columnSpecCnt == 2) + ? columnFamilySpec.getChild(3) + : null; + if (columnSpecCnt == 1) { - // table.cf['key']['column'] + assert columnTree != null; + + byte[] columnNameBytes = (columnTree.getType() == CliParser.FUNCTION_CALL) + ? convertValueByFunction(columnTree, null, null).array() + : columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), cfDef); + + if (isSuper) - superColumnName = columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), cfDef); + superColumnName = columnNameBytes; else - columnName = columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), cfDef); + columnName = columnNameBytes; } else if (columnSpecCnt == 2) { + assert columnTree != null; + assert subColumnTree != null; + // table.cf['key']['column']['column'] - superColumnName = columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), cfDef); - columnName = subColumnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 1), cfDef); + superColumnName = (columnTree.getType() == CliParser.FUNCTION_CALL) + ? convertValueByFunction(columnTree, null, null).array() + : columnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 0), cfDef); + + columnName = (subColumnTree.getType() == CliParser.FUNCTION_CALL) + ? convertValueByFunction(subColumnTree, null, null).array() + : subColumnNameAsByteArray(CliCompiler.getColumn(columnFamilySpec, 1), cfDef); } ColumnPath path = new ColumnPath(columnFamily); @@ -416,12 +443,11 @@ public class CliClient if (isCounterCF(cfDef)) { - thriftClient.remove_counter(ByteBufferUtil.bytes(key), path, consistencyLevel); + thriftClient.remove_counter(key, path, consistencyLevel); } else { - thriftClient.remove(ByteBufferUtil.bytes(key), path, - FBUtilities.timestampMicros(), consistencyLevel); + thriftClient.remove(key, path, FBUtilities.timestampMicros(), consistencyLevel); } sessionState.out.println(String.format("%s removed.", (columnSpecCnt == 0) ? "row" : "column")); } @@ -1050,11 +1076,15 @@ public class CliClient return; String cfName = CliCompiler.getColumnFamily(statement, keyspacesMap.get(keySpace).cf_defs); - // first child is a column family name - CfDef cfDef = getCfDef(cfName); try { + // request correct cfDef from the server + CfDef cfDef = getCfDef(thriftClient.describe_keyspace(this.keySpace), cfName); + + if (cfDef == null) + throw new RuntimeException("Column Family " + cfName + " was not found in the current keyspace."); + String mySchemaVersion = thriftClient.system_update_column_family(updateCfDefAttributes(statement, cfDef)); sessionState.out.println(mySchemaVersion); validateSchemaIsSettled(mySchemaVersion); @@ -1202,7 +1232,7 @@ public class CliClient cfDef.setKey_cache_save_period_in_seconds(Integer.parseInt(mValue)); break; case DEFAULT_VALIDATION_CLASS: - cfDef.setDefault_validation_class(mValue); + cfDef.setDefault_validation_class(CliUtils.unescapeSQLString(mValue)); break; case MIN_COMPACTION_THRESHOLD: cfDef.setMin_compaction_threshold(Integer.parseInt(mValue)); @@ -1252,6 +1282,9 @@ public class CliClient String version = thriftClient.system_drop_keyspace(keyspaceName); sessionState.out.println(version); validateSchemaIsSettled(version); + + if (keyspaceName.equals(keySpace)) //we just deleted the keyspace we were authenticated too + keySpace = null; } /** @@ -1898,7 +1931,18 @@ public class CliClient { return getCfDef(this.keySpace, columnFamilyName); } - + + private CfDef getCfDef(KsDef keyspace, String columnFamilyName) + { + for (CfDef cfDef : keyspace.cf_defs) + { + if (cfDef.name.equals(columnFamilyName)) + return cfDef; + } + + return null; + } + /** * Used to parse meta tree and compile meta attributes into List * @param cfDef - column family definition diff --git a/src/java/org/apache/cassandra/cli/CliUserHelp.java b/src/java/org/apache/cassandra/cli/CliUserHelp.java index a64bf8791e..5a9e231554 100644 --- a/src/java/org/apache/cassandra/cli/CliUserHelp.java +++ b/src/java/org/apache/cassandra/cli/CliUserHelp.java @@ -19,9 +19,6 @@ package org.apache.cassandra.cli; import java.util.List; -/** - * @author Pavel A. Yaskevich - */ public class CliUserHelp { public String banner; diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 01de9d2f5b..3bd728748a 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -101,7 +101,7 @@ public class DatabaseDescriptor try { url = new URL(configUrl); - url.openStream(); // catches well-formed but bogus URLs + url.openStream().close(); // catches well-formed but bogus URLs } catch (Exception e) { diff --git a/src/java/org/apache/cassandra/db/TruncateResponse.java b/src/java/org/apache/cassandra/db/TruncateResponse.java index 8929f3614c..3c1c507692 100644 --- a/src/java/org/apache/cassandra/db/TruncateResponse.java +++ b/src/java/org/apache/cassandra/db/TruncateResponse.java @@ -31,8 +31,6 @@ import org.apache.cassandra.utils.FBUtilities; /** * This message is sent back the truncate operation and basically specifies if * the truncate succeeded. - * - * @author rantav@gmail.com */ public class TruncateResponse { diff --git a/src/java/org/apache/cassandra/db/Truncation.java b/src/java/org/apache/cassandra/db/Truncation.java index 70b05e4b8e..fb2b10dedc 100644 --- a/src/java/org/apache/cassandra/db/Truncation.java +++ b/src/java/org/apache/cassandra/db/Truncation.java @@ -31,9 +31,6 @@ import org.apache.cassandra.utils.FBUtilities; /** * A truncate operation descriptor - * - * @author rantav@gmail.com - * */ public class Truncation implements MessageProducer { diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionController.java b/src/java/org/apache/cassandra/db/compaction/CompactionController.java index c7d5501c67..b724af60b9 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionController.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionController.java @@ -46,6 +46,7 @@ public class CompactionController public final boolean isMajor; public final int gcBefore; + private int throttleResolution; public CompactionController(ColumnFamilyStore cfs, Collection sstables, int gcBefore, boolean forceDeserialize) { @@ -55,15 +56,26 @@ public class CompactionController this.gcBefore = gcBefore; this.forceDeserialize = forceDeserialize; isMajor = cfs.isCompleteSSTables(this.sstables); + // how many rows we expect to compact in 100ms + long rowSize = cfs.getMeanRowSize(); + int rowsPerSecond = rowSize > 0 + ? (int) (DatabaseDescriptor.getCompactionThroughputMbPerSec() * 1024 * 1024 / rowSize) + : 1000; + throttleResolution = rowsPerSecond / 10; + if (throttleResolution <= 0) + throttleResolution = 1; + } + + public int getThrottleResolution() + { + return throttleResolution; } - /** @return the keyspace name */ public String getKeyspace() { return cfs.table.name; } - /** @return the column family name */ public String getColumnFamily() { return cfs.columnFamily; diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index ff6075f17c..55a52af2f3 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -121,9 +121,7 @@ implements CloseableIterator, CompactionInfo.Holder int newTarget = totalBytesPerMS / Math.max(1, CompactionManager.instance.getActiveCompactions()); if (newTarget != targetBytesPerMS) - logger.info(String.format("%s now compacting at %d bytes/ms.", - this, - newTarget)); + logger.debug("{} now compacting at {} bytes/ms.", this, newTarget); targetBytesPerMS = newTarget; // the excess bytes that were compacted in this period @@ -136,7 +134,14 @@ implements CloseableIterator, CompactionInfo.Holder if (logger.isTraceEnabled()) logger.trace(String.format("Compacted %d bytes in %d ms: throttling for %d ms", bytesSinceLast, msSinceLast, timeToDelay)); - try { Thread.sleep(timeToDelay); } catch (InterruptedException e) { throw new AssertionError(e); } + try + { + Thread.sleep(timeToDelay); + } + catch (InterruptedException e) + { + throw new AssertionError(e); + } } bytesAtLastDelay = bytesRead; timeAtLastDelay = System.currentTimeMillis(); diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index dbb3636106..a0bec4f70a 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -194,6 +194,11 @@ public class Gossiper implements IFailureDetectionEventListener versions.put(address, version); } + public void resetVersion(InetAddress endpoint) + { + versions.remove(endpoint); + } + public Integer getVersion(InetAddress address) { Integer v = versions.get(address); diff --git a/src/java/org/apache/cassandra/hadoop/ColumnFamilyInputFormat.java b/src/java/org/apache/cassandra/hadoop/ColumnFamilyInputFormat.java index 6415878274..51c8fcaf7f 100644 --- a/src/java/org/apache/cassandra/hadoop/ColumnFamilyInputFormat.java +++ b/src/java/org/apache/cassandra/hadoop/ColumnFamilyInputFormat.java @@ -35,8 +35,11 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.db.IColumn; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.Range; import org.apache.cassandra.thrift.Cassandra; import org.apache.cassandra.thrift.InvalidRequestException; +import org.apache.cassandra.thrift.KeyRange; import org.apache.cassandra.thrift.TokenRange; import org.apache.cassandra.thrift.TBinaryProtocol; import org.apache.hadoop.conf.Configuration; @@ -102,10 +105,44 @@ public class ColumnFamilyInputFormat extends InputFormat>> splitfutures = new ArrayList>>(); + KeyRange jobKeyRange = ConfigHelper.getInputKeyRange(conf); + IPartitioner partitioner = null; + Range jobRange = null; + if (jobKeyRange != null) + { + partitioner = ConfigHelper.getPartitioner(context.getConfiguration()); + assert partitioner.preservesOrder() : "ConfigHelper.setInputKeyRange(..) can only be used with a order preserving paritioner"; + assert jobKeyRange.start_key == null : "only start_token supported"; + assert jobKeyRange.end_key == null : "only end_token supported"; + jobRange = new Range(partitioner.getTokenFactory().fromString(jobKeyRange.start_token), + partitioner.getTokenFactory().fromString(jobKeyRange.end_token), + partitioner); + } + for (TokenRange range : masterRangeNodes) { + if (jobRange == null) + { // for each range, pick a live owner and ask it to compute bite-sized splits splitfutures.add(executor.submit(new SplitCallable(range, conf))); + } + else + { + Range dhtRange = new Range(partitioner.getTokenFactory().fromString(range.start_token), + partitioner.getTokenFactory().fromString(range.end_token), + partitioner); + + if (dhtRange.intersects(jobRange)) + { + Set intersections = dhtRange.intersectionWith(jobRange); + assert intersections.size() == 1 : "wrapping ranges not yet supported"; + Range intersection = intersections.iterator().next(); + range.start_token = partitioner.getTokenFactory().toString(intersection.left); + range.end_token = partitioner.getTokenFactory().toString(intersection.right); + // for each range, pick a live owner and ask it to compute bite-sized splits + splitfutures.add(executor.submit(new SplitCallable(range, conf))); + } + } } // wait until we have all the results back diff --git a/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordWriter.java b/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordWriter.java index 433c6efe59..3cd19ef63c 100644 --- a/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordWriter.java +++ b/src/java/org/apache/cassandra/hadoop/ColumnFamilyRecordWriter.java @@ -53,7 +53,6 @@ import org.apache.thrift.transport.TSocket; * directly to a responsible endpoint. *

* - * @author Karthick Sankarachary * @see ColumnFamilyOutputFormat * @see OutputFormat * diff --git a/src/java/org/apache/cassandra/hadoop/ConfigHelper.java b/src/java/org/apache/cassandra/hadoop/ConfigHelper.java index 0478ac7c12..8345fb7950 100644 --- a/src/java/org/apache/cassandra/hadoop/ConfigHelper.java +++ b/src/java/org/apache/cassandra/hadoop/ConfigHelper.java @@ -22,6 +22,7 @@ package org.apache.cassandra.hadoop; import org.apache.cassandra.config.ConfigurationException; import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.thrift.KeyRange; import org.apache.cassandra.thrift.SlicePredicate; import org.apache.cassandra.thrift.TBinaryProtocol; import org.apache.cassandra.utils.FBUtilities; @@ -42,6 +43,7 @@ public class ConfigHelper private static final String INPUT_COLUMNFAMILY_CONFIG = "cassandra.input.columnfamily"; private static final String OUTPUT_COLUMNFAMILY_CONFIG = "cassandra.output.columnfamily"; private static final String INPUT_PREDICATE_CONFIG = "cassandra.input.predicate"; + private static final String INPUT_KEYRANGE_CONFIG = "cassandra.input.keyRange"; private static final String OUTPUT_PREDICATE_CONFIG = "cassandra.output.predicate"; private static final String INPUT_SPLIT_SIZE_CONFIG = "cassandra.input.split.size"; private static final int DEFAULT_SPLIT_SIZE = 64 * 1024; @@ -195,6 +197,53 @@ public class ConfigHelper return predicate; } + /** + * Set the KeyRange to limit the rows. + * @param conf Job configuration you are about to run + */ + public static void setInputRange(Configuration conf, String startToken, String endToken) + { + KeyRange range = new KeyRange().setStart_token(startToken).setEnd_token(endToken); + conf.set(INPUT_KEYRANGE_CONFIG, keyRangeToString(range)); + } + + /** may be null if unset */ + public static KeyRange getInputKeyRange(Configuration conf) + { + String str = conf.get(INPUT_KEYRANGE_CONFIG); + return null != str ? keyRangeFromString(str) : null; + } + + private static String keyRangeToString(KeyRange keyRange) + { + assert keyRange != null; + TSerializer serializer = new TSerializer(new TBinaryProtocol.Factory()); + try + { + return FBUtilities.bytesToHex(serializer.serialize(keyRange)); + } + catch (TException e) + { + throw new RuntimeException(e); + } + } + + private static KeyRange keyRangeFromString(String st) + { + assert st != null; + TDeserializer deserializer = new TDeserializer(new TBinaryProtocol.Factory()); + KeyRange keyRange = new KeyRange(); + try + { + deserializer.deserialize(keyRange, FBUtilities.hexToBytes(st)); + } + catch (TException e) + { + throw new RuntimeException(e); + } + return keyRange; + } + public static String getInputKeyspace(Configuration conf) { return conf.get(INPUT_KEYSPACE_CONFIG); diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java b/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java index 6297f8f970..ea0a0cf0e7 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java @@ -101,7 +101,7 @@ public class SSTableLoader public LoaderFuture stream(Set toIgnore) throws IOException { - client.init(); + client.init(keyspace); Collection sstables = openSSTables(); if (sstables.isEmpty()) @@ -234,7 +234,7 @@ public class SSTableLoader * This method is guaranted to be called before any other method of a * client. */ - public abstract void init(); + public abstract void init(String keyspace); /** * Stop the client. diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java index 4e45043fa8..11eda6f13e 100644 --- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java +++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java @@ -87,9 +87,8 @@ public abstract class AbstractReplicationStrategy * we return a List to avoid an extra allocation when sorting by proximity later * @param searchToken the token the natural endpoints are requested for * @return a copy of the natural endpoints for the given token - * @throws IllegalStateException if the number of requested replicas is greater than the number of known endpoints */ - public ArrayList getNaturalEndpoints(Token searchToken) throws IllegalStateException + public ArrayList getNaturalEndpoints(Token searchToken) { Token keyToken = TokenMetadata.firstToken(tokenMetadata.sortedTokens(), searchToken); ArrayList endpoints = getCachedEndpoints(keyToken); @@ -99,10 +98,6 @@ public abstract class AbstractReplicationStrategy keyToken = TokenMetadata.firstToken(tokenMetadataClone.sortedTokens(), searchToken); endpoints = new ArrayList(calculateNaturalEndpoints(searchToken, tokenMetadataClone)); cacheEndpoint(keyToken, endpoints); - // calculateNaturalEndpoints should have checked this already, this is a safety - assert getReplicationFactor() <= endpoints.size() : String.format("endpoints %s generated for RF of %s", - Arrays.toString(endpoints.toArray()), - getReplicationFactor()); } return new ArrayList(endpoints); @@ -115,9 +110,8 @@ public abstract class AbstractReplicationStrategy * * @param searchToken the token the natural endpoints are requested for * @return a copy of the natural endpoints for the given token - * @throws IllegalStateException if the number of requested replicas is greater than the number of known endpoints */ - public abstract List calculateNaturalEndpoints(Token searchToken, TokenMetadata tokenMetadata) throws IllegalStateException; + public abstract List calculateNaturalEndpoints(Token searchToken, TokenMetadata tokenMetadata); public IWriteResponseHandler getWriteResponseHandler(Collection writeEndpoints, Multimap hintedEndpoints, diff --git a/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java b/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java index 317631a87d..0537ffbbcb 100644 --- a/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java +++ b/src/java/org/apache/cassandra/locator/NetworkTopologyStrategy.java @@ -120,9 +120,6 @@ public class NetworkTopologyStrategy extends AbstractReplicationStrategy dcEndpoints.add(endpoint); } - if (dcEndpoints.size() < dcReplicas) - throw new IllegalStateException(String.format("datacenter (%s) has no more endpoints, (%s) replicas still needed", - dcName, dcReplicas - dcEndpoints.size())); if (logger.isDebugEnabled()) logger.debug("{} endpoints in datacenter {} for token {} ", new Object[] { StringUtils.join(dcEndpoints, ","), dcName, searchToken}); diff --git a/src/java/org/apache/cassandra/locator/OldNetworkTopologyStrategy.java b/src/java/org/apache/cassandra/locator/OldNetworkTopologyStrategy.java index 558d6599ff..d8f32c889b 100644 --- a/src/java/org/apache/cassandra/locator/OldNetworkTopologyStrategy.java +++ b/src/java/org/apache/cassandra/locator/OldNetworkTopologyStrategy.java @@ -96,9 +96,6 @@ public class OldNetworkTopologyStrategy extends AbstractReplicationStrategy if (!endpoints.contains(metadata.getEndpoint(t))) endpoints.add(metadata.getEndpoint(t)); } - - if (endpoints.size() < replicas) - throw new IllegalStateException(String.format("replication factor (%s) exceeds number of endpoints (%s)", replicas, endpoints.size())); } return endpoints; diff --git a/src/java/org/apache/cassandra/locator/SimpleStrategy.java b/src/java/org/apache/cassandra/locator/SimpleStrategy.java index 09935df2cb..024e9d4e6d 100644 --- a/src/java/org/apache/cassandra/locator/SimpleStrategy.java +++ b/src/java/org/apache/cassandra/locator/SimpleStrategy.java @@ -56,10 +56,6 @@ public class SimpleStrategy extends AbstractReplicationStrategy { endpoints.add(metadata.getEndpoint(iter.next())); } - - if (endpoints.size() < replicas) - throw new IllegalStateException(String.format("replication factor (%s) exceeds number of endpoints (%s)", replicas, endpoints.size())); - return endpoints; } diff --git a/src/java/org/apache/cassandra/net/OutboundTcpConnection.java b/src/java/org/apache/cassandra/net/OutboundTcpConnection.java index fadd1b2069..828dac064e 100644 --- a/src/java/org/apache/cassandra/net/OutboundTcpConnection.java +++ b/src/java/org/apache/cassandra/net/OutboundTcpConnection.java @@ -30,6 +30,7 @@ import java.nio.ByteBuffer; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; +import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.utils.ByteBufferUtil; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -138,6 +139,9 @@ public class OutboundTcpConnection extends Thread output = null; socket = null; } + + // when we see the node again, try to connect at the most recent protocol we know about + Gossiper.instance.resetVersion(endpoint); } private ByteBuffer take() diff --git a/src/java/org/apache/cassandra/service/AbstractCassandraDaemon.java b/src/java/org/apache/cassandra/service/AbstractCassandraDaemon.java index 93c664a483..7ad7ba53b9 100644 --- a/src/java/org/apache/cassandra/service/AbstractCassandraDaemon.java +++ b/src/java/org/apache/cassandra/service/AbstractCassandraDaemon.java @@ -73,12 +73,30 @@ public abstract class AbstractCassandraDaemon implements CassandraDaemon } catch (MalformedURLException ex) { - // load from the classpath. + // then try loading from the classpath. configLocation = AbstractCassandraDaemon.class.getClassLoader().getResource(config); - if (configLocation == null) - throw new RuntimeException("Couldn't figure out log4j configuration."); } - PropertyConfigurator.configureAndWatch(configLocation.getFile(), 10000); + + if (configLocation == null) + throw new RuntimeException("Couldn't figure out log4j configuration: "+config); + + // Now convert URL to a filename + String configFileName = null; + try + { + // first try URL.getFile() which works for opaque URLs (file:foo) and paths without spaces + configFileName = configLocation.getFile(); + File configFile = new File(configFileName); + // then try alternative approach which works for all hierarchical URLs with or without spaces + if (!configFile.exists()) + configFileName = new File(configLocation.toURI()).getCanonicalPath(); + } + catch (Exception e) + { + throw new RuntimeException("Couldn't convert log4j configuration location to a valid file", e); + } + + PropertyConfigurator.configureAndWatch(configFileName, 10000); org.apache.log4j.Logger.getLogger(AbstractCassandraDaemon.class).info("Logging initialized"); } diff --git a/src/java/org/apache/cassandra/service/EmbeddedCassandraService.java b/src/java/org/apache/cassandra/service/EmbeddedCassandraService.java index 6ccf3e4645..53bcb16be8 100644 --- a/src/java/org/apache/cassandra/service/EmbeddedCassandraService.java +++ b/src/java/org/apache/cassandra/service/EmbeddedCassandraService.java @@ -45,8 +45,6 @@ import org.apache.thrift.transport.TTransportException; cassandra.start(); * - * @author Ran Tavory (rantav@gmail.com) - * */ public class EmbeddedCassandraService { diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 1370d1c0b8..8f6e19b803 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -1296,6 +1296,26 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe return stringify(Gossiper.instance.getUnreachableMembers()); } + public String[] getAllDataFileLocations() + { + return DatabaseDescriptor.getAllDataFileLocations(); + } + + public String[] getAllDataFileLocationsForTable(String table) + { + return DatabaseDescriptor.getAllDataFileLocationsForTable(table); + } + + public String getCommitLogLocation() + { + return DatabaseDescriptor.getCommitLogLocation(); + } + + public String getSavedCachesLocation() + { + return DatabaseDescriptor.getSavedCachesLocation(); + } + private List stringify(Iterable endpoints) { List stringEndpoints = new ArrayList(); @@ -2448,7 +2468,15 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe SSTableLoader.Client client = new SSTableLoader.Client() { - public void init() {} + public void init(String keyspace) + { + for (Map.Entry> entry : StorageService.instance.getRangeToAddressMap(keyspace).entrySet()) + { + Range range = entry.getKey(); + for (InetAddress endpoint : entry.getValue()) + addRangeForEndpoint(range, endpoint); + } + } public boolean validateColumnFamily(String keyspace, String cfName) { diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java index c3da511ac1..5451eb1022 100644 --- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java +++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java @@ -85,6 +85,31 @@ public interface StorageServiceMBean */ public String getReleaseVersion(); + /** + * Get the list of all data file locations from conf + * @return String array of all locations + */ + public String[] getAllDataFileLocations(); + + /** + * Get the list of data file locations for a given keyspace + * @param keyspace the keyspace to get locatiosn for. + * @return String array of all locations + */ + public String[] getAllDataFileLocationsForTable(String table); + + /** + * Get location of the commit log + * @return a string path + */ + public String getCommitLogLocation(); + + /** + * Get location of the saved caches dir + * @return a string path + */ + public String getSavedCachesLocation(); + /** * Retrieve a map of range to end points that describe the ring topology * of a Cassandra cluster. diff --git a/src/java/org/apache/cassandra/service/WriteResponseHandler.java b/src/java/org/apache/cassandra/service/WriteResponseHandler.java index ff6fb34969..435ce6c382 100644 --- a/src/java/org/apache/cassandra/service/WriteResponseHandler.java +++ b/src/java/org/apache/cassandra/service/WriteResponseHandler.java @@ -85,9 +85,9 @@ public class WriteResponseHandler extends AbstractWriteResponseHandler case THREE: return 3; case QUORUM: - return (writeEndpoints.size() / 2) + 1; + return (Table.open(table).getReplicationStrategy().getReplicationFactor() / 2) + 1; case ALL: - return writeEndpoints.size(); + return Table.open(table).getReplicationStrategy().getReplicationFactor(); default: throw new UnsupportedOperationException("invalid consistency level: " + consistencyLevel.toString()); } diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index 4baf84804d..e5979369cd 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -960,6 +960,7 @@ public class CassandraServer implements Cassandra.Iface CFMetaData oldCfm = DatabaseDescriptor.getCFMetaData(CFMetaData.getId(cf_def.keyspace, cf_def.name)); if (oldCfm == null) throw new InvalidRequestException("Could not find column family definition to modify."); + ThriftValidation.validateCfDef(cf_def, oldCfm); validateSchemaAgreement(); try diff --git a/src/java/org/apache/cassandra/tools/BulkLoader.java b/src/java/org/apache/cassandra/tools/BulkLoader.java index 45b7722ec8..36b7a6aae7 100644 --- a/src/java/org/apache/cassandra/tools/BulkLoader.java +++ b/src/java/org/apache/cassandra/tools/BulkLoader.java @@ -24,6 +24,7 @@ import java.net.InetAddress; import java.net.UnknownHostException; import java.util.*; +import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; @@ -57,7 +58,7 @@ public class BulkLoader LoaderOptions options = LoaderOptions.parseArgs(args); try { - SSTableLoader loader = new SSTableLoader(options.directory, new ExternalClient(options.directory.getName(), options), options); + SSTableLoader loader = new SSTableLoader(options.directory, new ExternalClient(options), options); SSTableLoader.LoaderFuture future = loader.stream(options.ignores); if (options.noProgress) @@ -164,18 +165,16 @@ public class BulkLoader static class ExternalClient extends SSTableLoader.Client { - private final String keyspace; private final Map> knownCfs = new HashMap>(); private final SSTableLoader.OutputHandler outputHandler; - public ExternalClient(String keyspace, SSTableLoader.OutputHandler outputHandler) + public ExternalClient(SSTableLoader.OutputHandler outputHandler) { super(); - this.keyspace = keyspace; this.outputHandler = outputHandler; } - public void init() + public void init(String keyspace) { outputHandler.output(String.format("Starting client (and waiting %d seconds for gossip) ...", StorageService.RING_DELAY / 1000)); try diff --git a/src/java/org/apache/cassandra/tools/NodeCmd.java b/src/java/org/apache/cassandra/tools/NodeCmd.java index 85a86c72b3..35792ba800 100644 --- a/src/java/org/apache/cassandra/tools/NodeCmd.java +++ b/src/java/org/apache/cassandra/tools/NodeCmd.java @@ -720,6 +720,15 @@ public class NodeCmd e.printStackTrace(); System.exit(3); } + + private static void complainNonzeroArgs(String[] args, NodeCommand cmd) + { + if (args.length > 0) { + System.err.println("Too many arguments for command '"+cmd.toString()+"'."); + printUsage(); + System.exit(1); + } + } private static void handleSnapshots(NodeCommand nc, String tag, String[] cmdArgs, NodeProbe probe) throws InterruptedException, IOException { diff --git a/src/resources/org/apache/cassandra/cli/CliHelp.yaml b/src/resources/org/apache/cassandra/cli/CliHelp.yaml index da15b7702a..9476920390 100644 --- a/src/resources/org/apache/cassandra/cli/CliHelp.yaml +++ b/src/resources/org/apache/cassandra/cli/CliHelp.yaml @@ -433,7 +433,7 @@ commands: store the whole values of its rows, so it is extremely space-intensive. It's best to only use the row cache if you have hot rows or static rows. - - keys_cache_save_period: Duration in seconds after which Cassandra should + - key_cache_save_period: Duration in seconds after which Cassandra should safe the keys cache. Caches are saved to saved_caches_directory as specified in conf/Cassandra.yaml. Default is 14400 or 4 hours. @@ -674,7 +674,7 @@ commands: store the whole values of its rows, so it is extremely space-intensive. It's best to only use the row cache if you have hot rows or static rows. - - keys_cache_save_period: Duration in seconds after which Cassandra should + - key_cache_save_period: Duration in seconds after which Cassandra should safe the keys cache. Caches are saved to saved_caches_directory as specified in conf/Cassandra.yaml. Default is 14400 or 4 hours. diff --git a/test/unit/org/apache/cassandra/cli/CliTest.java b/test/unit/org/apache/cassandra/cli/CliTest.java index 20836c5275..d630e0fbe9 100644 --- a/test/unit/org/apache/cassandra/cli/CliTest.java +++ b/test/unit/org/apache/cassandra/cli/CliTest.java @@ -54,6 +54,8 @@ public class CliTest extends CleanupHelper "get CF1 where world2 = long(15);", "get cF1 where world2 = long(15);", "get Cf1 where world2 = long(15);", + "del CF1[utf8('hello')][utf8('world')];", + "del CF1[hello][world2];", "set CF1['hello'][time_spent_uuid] = timeuuid(a8098c1a-f86e-11da-bd1a-00112444be1e);", "create column family CF2 with comparator=IntegerType;", "assume CF2 keys as utf8;", @@ -132,6 +134,10 @@ public class CliTest extends CleanupHelper "set sCf1['hello'][1][9999] = 938;", "set sCf1['hello'][1][9999] = 938 with ttl = 30;", "set sCf1['hello'][1][9999] = 938 with ttl = 560;", + "count sCf1[hello];", + "count sCf1[utf8('hello')];", + "count sCf1[utf8('hello')][integer(1)];", + "count sCF1[hello][1];", "list sCf1;", "del SCF1['hello'][1][9999];", "assume sCf1 comparator as utf8;", diff --git a/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java b/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java index 1c27f514f9..abeb57cdf4 100644 --- a/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java +++ b/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java @@ -36,9 +36,6 @@ import org.apache.cassandra.utils.ByteBufferUtil; /** * Test for the truncate operation. - * - * @author Ran Tavory (rantav@gmail.com) - * */ public class RecoveryManagerTruncateTest extends CleanupHelper {