From b85d44a25bec8355a86288376eca17021b9793f2 Mon Sep 17 00:00:00 2001 From: Vijay Parthasarathy Date: Thu, 29 Mar 2012 16:29:07 -0700 Subject: [PATCH 1/4] make ITC to handle versioning using BCA patch by Vijay; reviewed by Brandon Williams for CASSANDRA-4098 --- src/java/org/apache/cassandra/net/IncomingTcpConnection.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/java/org/apache/cassandra/net/IncomingTcpConnection.java b/src/java/org/apache/cassandra/net/IncomingTcpConnection.java index ee44a1c443..47ab39a663 100644 --- a/src/java/org/apache/cassandra/net/IncomingTcpConnection.java +++ b/src/java/org/apache/cassandra/net/IncomingTcpConnection.java @@ -48,7 +48,6 @@ public class IncomingTcpConnection extends Thread { assert socket != null; this.socket = socket; - from = socket.getInetAddress(); // maximize chance of this not being nulled by disconnect } /** @@ -94,6 +93,7 @@ public class IncomingTcpConnection extends Thread input = new DataInputStream(new BufferedInputStream(socket.getInputStream(), 4096)); // Receive the first message to set the version. Message msg = receiveMessage(input, version); + from = msg.getFrom(); // why? see => CASSANDRA-4099 if (version > MessagingService.version_) { // save the endpoint so gossip will reconnect to it @@ -102,7 +102,7 @@ public class IncomingTcpConnection extends Thread } else if (msg != null) { - Gossiper.instance.setVersion(msg.getFrom(), version); + Gossiper.instance.setVersion(from, version); logger.debug("set version for {} to {}", from, version); } From 7326ba88795665d241d2aac9a1386598f35f157e Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 30 Mar 2012 10:33:29 +0200 Subject: [PATCH 2/4] Allow custom types in CLI's assume command patch by xedin; reviewed by slebresne for CASSANDRA-4081 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/cli/Cli.g | 4 ++-- .../org/apache/cassandra/cli/CliClient.java | 20 +++++++++++++------ 3 files changed, 17 insertions(+), 8 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index e4d207c1d4..3316e87252 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -14,6 +14,7 @@ * ensure that directory is selected for compaction for user-defined tasks and upgradesstables (CASSANDRA-3985) * fix NPE on invalid CQL delete command (CASSANDRA-3755) + * allow custom types in CLI's assume command (CASSANDRA-4081) 1.0.8 diff --git a/src/java/org/apache/cassandra/cli/Cli.g b/src/java/org/apache/cassandra/cli/Cli.g index e7cba6c04c..742ccf2e1b 100644 --- a/src/java/org/apache/cassandra/cli/Cli.g +++ b/src/java/org/apache/cassandra/cli/Cli.g @@ -301,8 +301,8 @@ truncateStatement ; assumeStatement - : ASSUME columnFamily assumptionElement=Identifier 'AS' defaultType=Identifier - -> ^(NODE_ASSUME columnFamily $assumptionElement $defaultType) + : ASSUME columnFamily assumptionElement=Identifier 'AS' entityName + -> ^(NODE_ASSUME columnFamily $assumptionElement entityName) ; consistencyLevelStatement diff --git a/src/java/org/apache/cassandra/cli/CliClient.java b/src/java/org/apache/cassandra/cli/CliClient.java index 8e76b89022..dfbcb688aa 100644 --- a/src/java/org/apache/cassandra/cli/CliClient.java +++ b/src/java/org/apache/cassandra/cli/CliClient.java @@ -1491,17 +1491,25 @@ public class CliClient AbstractType comparator; // Could be UTF8Type, IntegerType, LexicalUUIDType etc. - String defaultType = statement.getChild(2).getText(); + String defaultType = CliUtils.unescapeSQLString(statement.getChild(2).getText()); try { - comparator = Function.valueOf(defaultType.toUpperCase()).getValidator(); + comparator = TypeParser.parse(defaultType); } - catch (Exception e) + catch (ConfigurationException e) { - String functions = Function.getFunctionNames(); - sessionState.out.println("Type '" + defaultType + "' was not found. Available: " + functions); - return; + try + { + comparator = Function.valueOf(defaultType.toUpperCase()).getValidator(); + } + catch (Exception ne) + { + String functions = Function.getFunctionNames(); + sessionState.out.println("Type '" + defaultType + "' was not found. Available: " + functions + + " Or any class which extends o.a.c.db.marshal.AbstractType."); + return; + } } // making string representation look property e.g. o.a.c.db.marshal.UTF8Type From 3931ee709da29d3b9d9c28b8d0ef34cfdb357c1c Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 30 Mar 2012 16:19:27 +0200 Subject: [PATCH 3/4] Fix total bytes count for parallel compaction patch by slebresne; reviewed by jbellis for CASSANDRA-3758 --- CHANGES.txt | 1 + .../db/compaction/AbstractCompactionIterable.java | 13 +++++++++++-- .../cassandra/db/compaction/CompactionIterable.java | 7 +------ .../db/compaction/ParallelCompactionIterable.java | 4 +--- 4 files changed, 14 insertions(+), 11 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 3316e87252..438bc91e02 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -15,6 +15,7 @@ tasks and upgradesstables (CASSANDRA-3985) * fix NPE on invalid CQL delete command (CASSANDRA-3755) * allow custom types in CLI's assume command (CASSANDRA-4081) + * Fix totalBytes count for parallel compactions (CASSANDRA-3758) 1.0.8 diff --git a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java index 53b1ba936f..31822192f5 100644 --- a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java +++ b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java @@ -41,15 +41,24 @@ public abstract class AbstractCompactionIterable implements Iterable scanners; protected final Throttle throttle; - public AbstractCompactionIterable(CompactionController controller, OperationType type) + public AbstractCompactionIterable(CompactionController controller, OperationType type, List scanners) { this.controller = controller; this.type = type; + this.scanners = scanners; + this.bytesRead = 0; + + long bytes = 0; + for (SSTableScanner scanner : scanners) + bytes += scanner.getFileLength(); + this.totalBytes = bytes; + this.throttle = new Throttle(toString(), new Throttle.ThroughputFunction() { /** @return Instantaneous throughput target in bytes per millisecond. */ diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterable.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterable.java index 5e0dfa71e7..65e4b54c42 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterable.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterable.java @@ -41,7 +41,6 @@ public class CompactionIterable extends AbstractCompactionIterable private static Logger logger = LoggerFactory.getLogger(CompactionIterable.class); private long row; - private final List scanners; private static final Comparator comparator = new Comparator() { @@ -58,12 +57,8 @@ public class CompactionIterable extends AbstractCompactionIterable protected CompactionIterable(OperationType type, List scanners, CompactionController controller) { - super(controller, type); - this.scanners = scanners; + super(controller, type, scanners); row = 0; - totalBytes = bytesRead = 0; - for (SSTableScanner scanner : scanners) - totalBytes += scanner.getFileLength(); } protected static List getScanners(Iterable sstables) throws IOException diff --git a/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java b/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java index dba8f558d8..52f81e0b19 100644 --- a/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java +++ b/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java @@ -59,7 +59,6 @@ public class ParallelCompactionIterable extends AbstractCompactionIterable { private static Logger logger = LoggerFactory.getLogger(ParallelCompactionIterable.class); - private final List scanners; private final int maxInMemorySize; public ParallelCompactionIterable(OperationType type, Iterable sstables, CompactionController controller) throws IOException @@ -74,8 +73,7 @@ public class ParallelCompactionIterable extends AbstractCompactionIterable protected ParallelCompactionIterable(OperationType type, List scanners, CompactionController controller, int maxInMemorySize) { - super(controller, type); - this.scanners = scanners; + super(controller, type, scanners); this.maxInMemorySize = maxInMemorySize; } From d69d304c4e4788daa201539e5e8441d30bb02396 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 30 Mar 2012 16:44:24 +0200 Subject: [PATCH 4/4] Remove dead sliced_buffer_size_in_kb option patch by slebresne; reviewed by jbellis for CASSANDRA-4076 --- CHANGES.txt | 1 + NEWS.txt | 3 ++- conf/cassandra.yaml | 4 ---- examples/client_only/conf/cassandra.yaml | 4 ---- src/java/org/apache/cassandra/config/Config.java | 2 -- .../apache/cassandra/config/DatabaseDescriptor.java | 10 ---------- .../db/columniterator/SSTableNamesIterator.java | 2 +- .../db/columniterator/SSTableSliceIterator.java | 2 +- .../org/apache/cassandra/io/sstable/SSTableReader.java | 2 +- .../apache/cassandra/io/sstable/SSTableReaderTest.java | 2 +- 10 files changed, 7 insertions(+), 25 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 55d17f755e..9c626798b2 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -2,6 +2,7 @@ * Adds caching and bloomFilterFpChange to CQL options (CASSANDRA-4042) * Adds posibility to autoconfigure size of the KeyCache (CASSANDRA-4087) * fix KEYS index from skipping results (CASSANDRA-3996) + * Remove sliced_buffer_size_in_kb dead option (CASSANDRA-4076) 1.1-beta2 diff --git a/NEWS.txt b/NEWS.txt index a588e91c7f..eca9e02c95 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -48,7 +48,8 @@ Upgrading + Prior to 1.1, you could use KEY as the primary key name in some select statements, even if the PK was actually given a different name. In 1.1+ you must use the defined PK name. - + - The sliced_buffer_size_in_kb option has been removed from the + cassandra.yaml config file (this option was a no-op since 1.0). Features -------- diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 866319cc19..1874f320b7 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -223,10 +223,6 @@ concurrent_writes: 32 # the maximum number of secondary indexes created on a single CF. memtable_flush_queue_size: 4 -# Buffer size to use when performing contiguous column slices. -# Increase this to the size of the column slices you typically perform -sliced_buffer_size_in_kb: 64 - # Whether to, when doing sequential writing, fsync() at intervals in # order to force the operating system to flush the dirty # buffers. Enable this to avoid sudden dirty buffer flushing from diff --git a/examples/client_only/conf/cassandra.yaml b/examples/client_only/conf/cassandra.yaml index 2d92794faf..701d1911df 100644 --- a/examples/client_only/conf/cassandra.yaml +++ b/examples/client_only/conf/cassandra.yaml @@ -150,10 +150,6 @@ concurrent_writes: 32 # By default this will be set to the amount of data directories defined. #memtable_flush_writers: 1 -# Buffer size to use when performing contiguous column slices. -# Increase this to the size of the column slices you typically perform -sliced_buffer_size_in_kb: 64 - # TCP port, for commands and data storage_port: 7000 diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index 2bc34dce77..91b96f16d8 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -61,8 +61,6 @@ public class Config public Integer memtable_flush_writers = null; // will get set to the length of data dirs in DatabaseDescriptor public Integer memtable_total_space_in_mb; - public Integer sliced_buffer_size_in_kb = 64; - public Integer storage_port = 7000; public Integer ssl_storage_port = 7001; public String listen_address; diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 8c5997f658..792dc186a1 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -864,16 +864,6 @@ public class DatabaseDescriptor return indexAccessMode; } - public static int getIndexedReadBufferSizeInKB() - { - return conf.column_index_size_in_kb; - } - - public static int getSlicedReadBufferSizeInKB() - { - return conf.sliced_buffer_size_in_kb; - } - public static boolean isSnapshotBeforeCompaction() { return conf.snapshot_before_compaction; diff --git a/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java b/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java index fdb0c5c085..c063d3e994 100644 --- a/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java +++ b/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java @@ -58,7 +58,7 @@ public class SSTableNamesIterator extends SimpleAbstractColumnIterator implement this.columns = columns; this.key = key; - FileDataInput file = sstable.getFileDataInput(key, DatabaseDescriptor.getIndexedReadBufferSizeInKB() * 1024); + FileDataInput file = sstable.getFileDataInput(key); if (file == null) return; diff --git a/src/java/org/apache/cassandra/db/columniterator/SSTableSliceIterator.java b/src/java/org/apache/cassandra/db/columniterator/SSTableSliceIterator.java index 72fe1bf3ee..8a43e28e41 100644 --- a/src/java/org/apache/cassandra/db/columniterator/SSTableSliceIterator.java +++ b/src/java/org/apache/cassandra/db/columniterator/SSTableSliceIterator.java @@ -45,7 +45,7 @@ public class SSTableSliceIterator implements IColumnIterator public SSTableSliceIterator(SSTableReader sstable, DecoratedKey key, ByteBuffer startColumn, ByteBuffer finishColumn, boolean reversed) { this.key = key; - fileToClose = sstable.getFileDataInput(this.key, DatabaseDescriptor.getSlicedReadBufferSizeInKB() * 1024); + fileToClose = sstable.getFileDataInput(this.key); if (fileToClose == null) return; diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java index 57cdfa2a96..6c679ee128 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java @@ -834,7 +834,7 @@ public class SSTableReader extends SSTable return new SSTableBoundedScanner(this, true, range); } - public FileDataInput getFileDataInput(DecoratedKey decoratedKey, int bufferSize) + public FileDataInput getFileDataInput(DecoratedKey decoratedKey) { long position = getPosition(decoratedKey, Operator.EQ); if (position < 0) diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java index 102fccd45b..004dac812c 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java @@ -125,7 +125,7 @@ public class SSTableReaderTest extends SchemaLoader for (int j = 0; j < 100; j += 2) { DecoratedKey dk = Util.dk(String.valueOf(j)); - FileDataInput file = sstable.getFileDataInput(dk, DatabaseDescriptor.getIndexedReadBufferSizeInKB() * 1024); + FileDataInput file = sstable.getFileDataInput(dk); DecoratedKey keyInDisk = SSTableReader.decodeKey(sstable.partitioner, sstable.descriptor, ByteBufferUtil.readWithShortLength(file));