From 6ae8adaa3c9eef3606ef1f6a27b828066225aac2 Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Thu, 29 Jan 2015 09:41:49 +0000 Subject: [PATCH 1/2] ninja follow-up to 8619 --- .../io/sstable/CQLSSTableWriter.java | 2 +- .../sstable/SSTableSimpleUnsortedWriter.java | 32 +++++++++++++------ 2 files changed, 23 insertions(+), 11 deletions(-) diff --git a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java index d58b28f487..8006112902 100644 --- a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java @@ -534,7 +534,7 @@ public class CQLSSTableWriter implements Closeable }; } - protected void addColumn(Cell cell) throws IOException + protected void addColumn(Column column) throws IOException { throw new UnsupportedOperationException(); } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java index b58e5741ba..db03ea17a1 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableSimpleUnsortedWriter.java @@ -23,6 +23,7 @@ import java.util.Map; import java.util.TreeMap; import java.util.concurrent.BlockingQueue; import java.util.concurrent.SynchronousQueue; +import java.util.concurrent.TimeUnit; import com.google.common.base.Throwables; @@ -165,17 +166,21 @@ public class SSTableSimpleUnsortedWriter extends AbstractSSTableSimpleWriter if (buffer.isEmpty()) return; - checkForWriterException(); - - columnFamily = null; - try + while (true) { - writeQueue.put(buffer); - } - catch (InterruptedException e) - { - throw new RuntimeException(e); + checkForWriterException(); + columnFamily = null; + try + { + if (writeQueue.offer(buffer, 1L, TimeUnit.SECONDS)) + break; + } + catch (InterruptedException e) + { + throw new RuntimeException(e); + + } } buffer = new Buffer(); currentSize = 0; @@ -213,8 +218,15 @@ public class SSTableSimpleUnsortedWriter extends AbstractSSTableSimpleWriter return; writer = getWriter(); + boolean first = true; for (Map.Entry entry : b.entrySet()) - writer.append(entry.getKey(), entry.getValue()); + { + if (entry.getValue().getColumnCount() > 0) + writer.append(entry.getKey(), entry.getValue()); + else if (!first) + throw new AssertionError("Empty partition"); + first = false; + } writer.close(); } } From d0005a83ef9bfc1042747cbe45708ded1f2002ba Mon Sep 17 00:00:00 2001 From: Jeff Jirsa Date: Fri, 30 Jan 2015 16:33:56 -0600 Subject: [PATCH 2/2] Disallow default_time_to_live on counter tables Patch by Jeff Jirsa; reviewed by Tyler Hobbs for CASSANDRA-8678 --- CHANGES.txt | 5 ++++- src/java/org/apache/cassandra/config/CFMetaData.java | 5 +++++ .../cassandra/cql3/statements/AlterTableStatement.java | 4 ++++ .../org/apache/cassandra/cql3/statements/CFPropDefs.java | 5 +++++ .../cassandra/cql3/statements/CreateTableStatement.java | 6 ++++++ src/java/org/apache/cassandra/db/marshal/AbstractType.java | 5 +++++ .../org/apache/cassandra/db/marshal/CounterColumnType.java | 5 +++++ 7 files changed, 34 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index 7fa5f63d74..4754867a73 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,8 @@ 2.0.13: - * Fix SSTableSimpleUnsortedWriter ConcurrentModificationException (CASSANDRA-8619) + * Prevent non-zero default_time_to_live on tables with counters + (CASSANDRA-8678) + * Fix SSTableSimpleUnsortedWriter ConcurrentModificationException + (CASSANDRA-8619) * Round up time deltas lower than 1ms in BulkLoader (CASSANDRA-8645) * Add batch remove iterator to ABSC (CASSANDRA-8414, 8666) diff --git a/src/java/org/apache/cassandra/config/CFMetaData.java b/src/java/org/apache/cassandra/config/CFMetaData.java index 9d69710b58..04a5b01b1e 100644 --- a/src/java/org/apache/cassandra/config/CFMetaData.java +++ b/src/java/org/apache/cassandra/config/CFMetaData.java @@ -2192,6 +2192,11 @@ public final class CFMetaData return true; } + public boolean isCounter() + { + return defaultValidator.isCounter(); + } + public boolean hasStaticColumns() { return !staticColumns.isEmpty(); diff --git a/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java b/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java index 32f949f2d7..f74670f14d 100644 --- a/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/AlterTableStatement.java @@ -268,6 +268,10 @@ public class AlterTableStatement extends SchemaAlteringStatement throw new InvalidRequestException(String.format("ALTER COLUMNFAMILY WITH invoked, but no parameters found")); cfProps.validate(); + + if (meta.isCounter() && cfProps.getDefaultTimeToLive() > 0) + throw new InvalidRequestException("Cannot set default_time_to_live on a table with counters"); + cfProps.applyToCFMetadata(cfm); break; case RENAME: diff --git a/src/java/org/apache/cassandra/cql3/statements/CFPropDefs.java b/src/java/org/apache/cassandra/cql3/statements/CFPropDefs.java index 6ce640645f..11aa9b2278 100644 --- a/src/java/org/apache/cassandra/cql3/statements/CFPropDefs.java +++ b/src/java/org/apache/cassandra/cql3/statements/CFPropDefs.java @@ -138,6 +138,11 @@ public class CFPropDefs extends PropertyDefinitions return compressionOptions; } + public Integer getDefaultTimeToLive() throws SyntaxException + { + return getInt(KW_DEFAULT_TIME_TO_LIVE, 0); + } + public void applyToCFMetadata(CFMetaData cfm) throws ConfigurationException, SyntaxException { if (hasProperty(KW_COMMENT)) diff --git a/src/java/org/apache/cassandra/cql3/statements/CreateTableStatement.java b/src/java/org/apache/cassandra/cql3/statements/CreateTableStatement.java index 5160d6f5f4..5215118a73 100644 --- a/src/java/org/apache/cassandra/cql3/statements/CreateTableStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/CreateTableStatement.java @@ -207,6 +207,7 @@ public class CreateTableStatement extends SchemaAlteringStatement CreateTableStatement stmt = new CreateTableStatement(cfName, properties, ifNotExists, staticColumns); + boolean hasCounters = false; Map definedCollections = null; for (Map.Entry entry : definitions.entrySet()) { @@ -218,6 +219,9 @@ public class CreateTableStatement extends SchemaAlteringStatement definedCollections = new HashMap(); definedCollections.put(id.key, (CollectionType)pt.getType()); } + else if (entry.getValue().isCounter()) + hasCounters = true; + stmt.columns.put(id, pt.getType()); // we'll remove what is not a column below } @@ -225,6 +229,8 @@ public class CreateTableStatement extends SchemaAlteringStatement throw new InvalidRequestException("No PRIMARY KEY specifed (exactly one required)"); else if (keyAliases.size() > 1) throw new InvalidRequestException("Multiple PRIMARY KEYs specifed (exactly one required)"); + else if (hasCounters && properties.getDefaultTimeToLive() > 0) + throw new InvalidRequestException("Cannot set default_time_to_live on a table with counters"); List kAliases = keyAliases.get(0); diff --git a/src/java/org/apache/cassandra/db/marshal/AbstractType.java b/src/java/org/apache/cassandra/db/marshal/AbstractType.java index e92f272421..dce521b49c 100644 --- a/src/java/org/apache/cassandra/db/marshal/AbstractType.java +++ b/src/java/org/apache/cassandra/db/marshal/AbstractType.java @@ -212,6 +212,11 @@ public abstract class AbstractType implements Comparator return false; } + public boolean isCounter() + { + return false; + } + public static AbstractType parseDefaultParameters(AbstractType baseType, TypeParser parser) throws SyntaxException { Map parameters = parser.getKeyValueParameters(); diff --git a/src/java/org/apache/cassandra/db/marshal/CounterColumnType.java b/src/java/org/apache/cassandra/db/marshal/CounterColumnType.java index 6a774580fa..1ecdad9ed4 100644 --- a/src/java/org/apache/cassandra/db/marshal/CounterColumnType.java +++ b/src/java/org/apache/cassandra/db/marshal/CounterColumnType.java @@ -31,6 +31,11 @@ public class CounterColumnType extends AbstractCommutativeType CounterColumnType() {} // singleton + public boolean isCounter() + { + return true; + } + public int compare(ByteBuffer o1, ByteBuffer o2) { if (o1 == null)