From d867ac1f41c59b31f8fb4f54a06c0118018cfc81 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 24 Sep 2015 13:24:57 -0700 Subject: [PATCH 1/3] Simplify row cache invalidation code patch by Jonathan Ellis; reviewed by Aleksey Yeschenko for CASSANDRA-10396 --- .../apache/cassandra/db/ColumnFamilyStore.java | 15 +++------------ .../db/compaction/CompactionController.java | 4 +--- .../db/compaction/CompactionManager.java | 2 +- .../cassandra/io/sstable/SSTableRewriter.java | 3 +-- .../apache/cassandra/streaming/StreamReader.java | 2 +- .../org/apache/cassandra/db/RowCacheTest.java | 2 +- .../db/compaction/CompactionsPurgeTest.java | 2 +- 7 files changed, 9 insertions(+), 21 deletions(-) diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index d553f4db58..b112e0e449 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -1226,7 +1226,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return String.format("%.2f/%.2f", onHeap, offHeap); } - public void maybeUpdateRowCache(DecoratedKey key) + public void maybeInvalidateCachedRow(DecoratedKey key) { if (!isRowCacheEnabled()) return; @@ -1247,7 +1247,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean long start = System.nanoTime(); Memtable mt = data.getMemtableFor(opGroup, replayPosition); final long timeDelta = mt.put(key, columnFamily, indexer, opGroup); - maybeUpdateRowCache(key); + maybeInvalidateCachedRow(key); metric.samplers.get(Sampler.WRITES).addSample(key.getKey(), key.hashCode(), 1); metric.writeLatency.addNano(System.nanoTime() - start); if(timeDelta < Long.MAX_VALUE) @@ -2047,7 +2047,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean RowCacheKey key = keyIter.next(); DecoratedKey dk = partitioner.decorateKey(ByteBuffer.wrap(key.key)); if (key.ksAndCFName.equals(metadata.ksAndCFName) && !Range.isInRanges(dk.getToken(), ranges)) - invalidateCachedRow(dk); + maybeInvalidateCachedRow(dk); } if (metadata.isCounter()) @@ -2532,15 +2532,6 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean CacheService.instance.rowCache.remove(key); } - public void invalidateCachedRow(DecoratedKey key) - { - UUID cfId = Schema.instance.getId(keyspace.getName(), this.name); - if (cfId == null) - return; // secondary index - - invalidateCachedRow(new RowCacheKey(metadata.ksAndCFName, key)); - } - public ClockAndCount getCachedCounter(ByteBuffer partitionKey, CellName cellName) { if (CacheService.instance.counterCache.getCapacity() == 0L) // counter cache disabled. diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionController.java b/src/java/org/apache/cassandra/db/compaction/CompactionController.java index 5f0a198689..24ef8433db 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionController.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionController.java @@ -24,8 +24,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.db.ColumnFamilyStore; -import org.apache.cassandra.db.lifecycle.SSTableIntervalTree; -import org.apache.cassandra.db.lifecycle.Tracker; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.RowPosition; import org.apache.cassandra.utils.AlwaysPresentFilter; @@ -191,7 +189,7 @@ public class CompactionController implements AutoCloseable public void invalidateCachedRow(DecoratedKey key) { - cfs.invalidateCachedRow(key); + cfs.maybeInvalidateCachedRow(key); } public void close() diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 0c6e24f002..8537aca0cd 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -934,7 +934,7 @@ public class CompactionManager implements CompactionManagerMBean if (Range.isInRanges(row.getKey().getToken(), ranges)) return row; - cfs.invalidateCachedRow(row.getKey()); + cfs.maybeInvalidateCachedRow(row.getKey()); if (indexedColumnsInRow != null) indexedColumnsInRow.clear(); diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java index dc4fe75bd8..b08b0388c0 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java @@ -20,7 +20,6 @@ package org.apache.cassandra.io.sstable; import java.util.*; import com.google.common.annotations.VisibleForTesting; -import com.google.common.util.concurrent.Runnables; import org.apache.cassandra.cache.InstrumentingCache; import org.apache.cassandra.cache.KeyCacheKey; @@ -119,7 +118,7 @@ public class SSTableRewriter extends Transactional.AbstractTransactional impleme { if (index == null) { - cfs.invalidateCachedRow(row.key); + cfs.maybeInvalidateCachedRow(row.key); } else { diff --git a/src/java/org/apache/cassandra/streaming/StreamReader.java b/src/java/org/apache/cassandra/streaming/StreamReader.java index 1ccebb0901..88591da051 100644 --- a/src/java/org/apache/cassandra/streaming/StreamReader.java +++ b/src/java/org/apache/cassandra/streaming/StreamReader.java @@ -171,6 +171,6 @@ public class StreamReader { DecoratedKey key = StorageService.getPartitioner().decorateKey(ByteBufferUtil.readWithShortLength(in)); writer.appendFromStream(key, cfs.metadata, in, inputVersion); - cfs.invalidateCachedRow(key); + cfs.maybeInvalidateCachedRow(key); } } diff --git a/test/unit/org/apache/cassandra/db/RowCacheTest.java b/test/unit/org/apache/cassandra/db/RowCacheTest.java index 5912d7c3bd..69c831dbf2 100644 --- a/test/unit/org/apache/cassandra/db/RowCacheTest.java +++ b/test/unit/org/apache/cassandra/db/RowCacheTest.java @@ -133,7 +133,7 @@ public class RowCacheTest int keysLeft = 109; for (int i = 109; i >= 10; i--) { - cachedStore.invalidateCachedRow(Util.dk("key" + i)); + cachedStore.maybeInvalidateCachedRow(Util.dk("key" + i)); assert CacheService.instance.rowCache.size() == keysLeft; keysLeft--; } diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java index e5baab602f..bfe8042b07 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java @@ -127,7 +127,7 @@ public class CompactionsPurgeTest // major compact and test that all columns but the resurrected one is completely gone FBUtilities.waitOnFutures(CompactionManager.instance.submitMaximal(cfs, Integer.MAX_VALUE, false)); - cfs.invalidateCachedRow(key); + cfs.maybeInvalidateCachedRow(key); ColumnFamily cf = cfs.getColumnFamily(QueryFilter.getIdentityFilter(key, cfName, System.currentTimeMillis())); assertColumns(cf, "5"); assertNotNull(cf.getColumn(cellname(String.valueOf(5)))); From 48b685e8521ea54d93c0d8d9e4ea80ecb1400dce Mon Sep 17 00:00:00 2001 From: Aleksey Yeschenko Date: Mon, 9 Nov 2015 20:15:44 +0000 Subject: [PATCH 2/3] Revert "Simplify row cache invalidation code" This reverts commit d867ac1f41c59b31f8fb4f54a06c0118018cfc81. --- .../apache/cassandra/db/ColumnFamilyStore.java | 15 ++++++++++++--- .../db/compaction/CompactionController.java | 4 +++- .../db/compaction/CompactionManager.java | 2 +- .../cassandra/io/sstable/SSTableRewriter.java | 3 ++- .../apache/cassandra/streaming/StreamReader.java | 2 +- .../org/apache/cassandra/db/RowCacheTest.java | 2 +- .../db/compaction/CompactionsPurgeTest.java | 2 +- 7 files changed, 21 insertions(+), 9 deletions(-) diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index b112e0e449..d553f4db58 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -1226,7 +1226,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return String.format("%.2f/%.2f", onHeap, offHeap); } - public void maybeInvalidateCachedRow(DecoratedKey key) + public void maybeUpdateRowCache(DecoratedKey key) { if (!isRowCacheEnabled()) return; @@ -1247,7 +1247,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean long start = System.nanoTime(); Memtable mt = data.getMemtableFor(opGroup, replayPosition); final long timeDelta = mt.put(key, columnFamily, indexer, opGroup); - maybeInvalidateCachedRow(key); + maybeUpdateRowCache(key); metric.samplers.get(Sampler.WRITES).addSample(key.getKey(), key.hashCode(), 1); metric.writeLatency.addNano(System.nanoTime() - start); if(timeDelta < Long.MAX_VALUE) @@ -2047,7 +2047,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean RowCacheKey key = keyIter.next(); DecoratedKey dk = partitioner.decorateKey(ByteBuffer.wrap(key.key)); if (key.ksAndCFName.equals(metadata.ksAndCFName) && !Range.isInRanges(dk.getToken(), ranges)) - maybeInvalidateCachedRow(dk); + invalidateCachedRow(dk); } if (metadata.isCounter()) @@ -2532,6 +2532,15 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean CacheService.instance.rowCache.remove(key); } + public void invalidateCachedRow(DecoratedKey key) + { + UUID cfId = Schema.instance.getId(keyspace.getName(), this.name); + if (cfId == null) + return; // secondary index + + invalidateCachedRow(new RowCacheKey(metadata.ksAndCFName, key)); + } + public ClockAndCount getCachedCounter(ByteBuffer partitionKey, CellName cellName) { if (CacheService.instance.counterCache.getCapacity() == 0L) // counter cache disabled. diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionController.java b/src/java/org/apache/cassandra/db/compaction/CompactionController.java index 24ef8433db..5f0a198689 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionController.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionController.java @@ -24,6 +24,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.lifecycle.SSTableIntervalTree; +import org.apache.cassandra.db.lifecycle.Tracker; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.RowPosition; import org.apache.cassandra.utils.AlwaysPresentFilter; @@ -189,7 +191,7 @@ public class CompactionController implements AutoCloseable public void invalidateCachedRow(DecoratedKey key) { - cfs.maybeInvalidateCachedRow(key); + cfs.invalidateCachedRow(key); } public void close() diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 8537aca0cd..0c6e24f002 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -934,7 +934,7 @@ public class CompactionManager implements CompactionManagerMBean if (Range.isInRanges(row.getKey().getToken(), ranges)) return row; - cfs.maybeInvalidateCachedRow(row.getKey()); + cfs.invalidateCachedRow(row.getKey()); if (indexedColumnsInRow != null) indexedColumnsInRow.clear(); diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java index b08b0388c0..dc4fe75bd8 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableRewriter.java @@ -20,6 +20,7 @@ package org.apache.cassandra.io.sstable; import java.util.*; import com.google.common.annotations.VisibleForTesting; +import com.google.common.util.concurrent.Runnables; import org.apache.cassandra.cache.InstrumentingCache; import org.apache.cassandra.cache.KeyCacheKey; @@ -118,7 +119,7 @@ public class SSTableRewriter extends Transactional.AbstractTransactional impleme { if (index == null) { - cfs.maybeInvalidateCachedRow(row.key); + cfs.invalidateCachedRow(row.key); } else { diff --git a/src/java/org/apache/cassandra/streaming/StreamReader.java b/src/java/org/apache/cassandra/streaming/StreamReader.java index 88591da051..1ccebb0901 100644 --- a/src/java/org/apache/cassandra/streaming/StreamReader.java +++ b/src/java/org/apache/cassandra/streaming/StreamReader.java @@ -171,6 +171,6 @@ public class StreamReader { DecoratedKey key = StorageService.getPartitioner().decorateKey(ByteBufferUtil.readWithShortLength(in)); writer.appendFromStream(key, cfs.metadata, in, inputVersion); - cfs.maybeInvalidateCachedRow(key); + cfs.invalidateCachedRow(key); } } diff --git a/test/unit/org/apache/cassandra/db/RowCacheTest.java b/test/unit/org/apache/cassandra/db/RowCacheTest.java index 69c831dbf2..5912d7c3bd 100644 --- a/test/unit/org/apache/cassandra/db/RowCacheTest.java +++ b/test/unit/org/apache/cassandra/db/RowCacheTest.java @@ -133,7 +133,7 @@ public class RowCacheTest int keysLeft = 109; for (int i = 109; i >= 10; i--) { - cachedStore.maybeInvalidateCachedRow(Util.dk("key" + i)); + cachedStore.invalidateCachedRow(Util.dk("key" + i)); assert CacheService.instance.rowCache.size() == keysLeft; keysLeft--; } diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java index bfe8042b07..e5baab602f 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java @@ -127,7 +127,7 @@ public class CompactionsPurgeTest // major compact and test that all columns but the resurrected one is completely gone FBUtilities.waitOnFutures(CompactionManager.instance.submitMaximal(cfs, Integer.MAX_VALUE, false)); - cfs.maybeInvalidateCachedRow(key); + cfs.invalidateCachedRow(key); ColumnFamily cf = cfs.getColumnFamily(QueryFilter.getIdentityFilter(key, cfName, System.currentTimeMillis())); assertColumns(cf, "5"); assertNotNull(cf.getColumn(cellname(String.valueOf(5)))); From 3674ad9dab8f29173d7d4ee82488a8e9ea586240 Mon Sep 17 00:00:00 2001 From: Paulo Motta Date: Thu, 22 Oct 2015 11:38:31 -0700 Subject: [PATCH 3/3] Reject counter writes in CQLSSTableWriter patch by Paulo Motta; reviewed by Aleksey Yeschenko for CASSANDRA-10258 --- CHANGES.txt | 1 + .../io/sstable/CQLSSTableWriter.java | 2 ++ .../io/sstable/CQLSSTableWriterTest.java | 22 +++++++++++++++++++ 3 files changed, 25 insertions(+) diff --git a/CHANGES.txt b/CHANGES.txt index 123c1f38b0..fa2017a152 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.1.12 + * Reject counter writes in CQLSSTableWriter (CASSANDRA-10258) * Remove superfluous COUNTER_MUTATION stage mapping (CASSANDRA-10605) * Improve json2sstable error reporting on nonexistent columns (CASSANDRA-10401) * (cqlsh) fix COPY using wrong variable name for time_format (CASSANDRA-10633) diff --git a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java index c364171126..ae8a3920bf 100644 --- a/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/CQLSSTableWriter.java @@ -453,6 +453,8 @@ public class CQLSSTableWriter implements Closeable this.boundNames = p.right; if (this.insert.hasConditions()) throw new IllegalArgumentException("Conditional statements are not supported"); + if (this.insert.isCounter()) + throw new IllegalArgumentException("Counter update statements are not supported"); if (this.boundNames.isEmpty()) throw new IllegalArgumentException("Provided insert statement has no bind variables"); return this; diff --git a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java index fa5cbb4b73..9c8a2c2de0 100644 --- a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java @@ -135,6 +135,28 @@ public class CQLSSTableWriterTest assertEquals(12, row.getInt("v2")); } + @Test(expected = IllegalArgumentException.class) + public void testForbidCounterUpdates() throws Exception + { + String KS = "cql_keyspace"; + String TABLE = "counter1"; + + File tempdir = Files.createTempDir(); + File dataDir = new File(tempdir.getAbsolutePath() + File.separator + KS + File.separator + TABLE); + assert dataDir.mkdirs(); + + String schema = "CREATE TABLE cql_keyspace.counter1 (" + + " my_id int, " + + " my_counter counter, " + + " PRIMARY KEY (my_id)" + + ")"; + String insert = String.format("UPDATE cql_keyspace.counter1 SET my_counter = my_counter - ? WHERE my_id = ?"); + CQLSSTableWriter.builder().inDirectory(dataDir) + .forTable(schema) + .withPartitioner(StorageService.instance.getPartitioner()) + .using(insert).build(); + } + @Test public void testSyncWithinPartition() throws Exception {