From 9a62ef339c2fc7f25cde102a052899438ef08927 Mon Sep 17 00:00:00 2001 From: Yuki Morishita Date: Fri, 28 Feb 2014 15:01:29 -0600 Subject: [PATCH 1/7] Replace differencers set with AtomicInteger to track sync complete --- .../apache/cassandra/repair/RepairJob.java | 21 +++++++++---------- .../cassandra/repair/RepairSession.java | 2 +- 2 files changed, 11 insertions(+), 12 deletions(-) diff --git a/src/java/org/apache/cassandra/repair/RepairJob.java b/src/java/org/apache/cassandra/repair/RepairJob.java index 475d7f7714..13fe51174c 100644 --- a/src/java/org/apache/cassandra/repair/RepairJob.java +++ b/src/java/org/apache/cassandra/repair/RepairJob.java @@ -19,14 +19,13 @@ package org.apache.cassandra.repair; import java.net.InetAddress; import java.util.*; -import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.Condition; import com.google.common.util.concurrent.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.concurrent.NamedThreadFactory; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; @@ -51,12 +50,13 @@ public class RepairJob private final List trees = new ArrayList<>(); // once all responses are received, each tree is compared with each other, and differencer tasks // are submitted. the job is done when all differencers are complete. - private final Set differencers = new HashSet<>(); private final ListeningExecutorService taskExecutor; private final Condition requestsSent = new SimpleCondition(); private int gcBefore = -1; private volatile boolean failed = false; + /* Count down as sync completes */ + private AtomicInteger waitForSync; /** * Create repair job to run on specific columnfamily @@ -172,7 +172,7 @@ public class RepairJob public void submitDifferencers() { assert !failed; - + List differencers = new ArrayList<>(); // We need to difference all trees one against another for (int i = 0; i < trees.size() - 1; ++i) { @@ -183,21 +183,20 @@ public class RepairJob Differencer differencer = new Differencer(desc, r1, r2); differencers.add(differencer); logger.debug("Queueing comparison {}", differencer); - taskExecutor.submit(differencer); } } + waitForSync = new AtomicInteger(differencers.size()); + for (Differencer differencer : differencers) + taskExecutor.submit(differencer); + trees.clear(); // allows gc to do its thing } /** * @return true if the given node pair was the last remaining */ - synchronized boolean completedSynchronization(NodePair nodes, boolean success) + boolean completedSynchronization() { - if (!success) - failed = true; - Differencer completed = new Differencer(desc, new TreeResponse(nodes.endpoint1, null), new TreeResponse(nodes.endpoint2, null)); - differencers.remove(completed); - return differencers.size() == 0; + return waitForSync.decrementAndGet() == 0; } } diff --git a/src/java/org/apache/cassandra/repair/RepairSession.java b/src/java/org/apache/cassandra/repair/RepairSession.java index 7ffe87ff48..ea31ff3702 100644 --- a/src/java/org/apache/cassandra/repair/RepairSession.java +++ b/src/java/org/apache/cassandra/repair/RepairSession.java @@ -211,7 +211,7 @@ public class RepairSession extends WrappedRunnable implements IEndpointStateChan logger.debug(String.format("[repair #%s] Repair completed between %s and %s on %s", getId(), nodes.endpoint1, nodes.endpoint2, desc.columnFamily)); - if (job.completedSynchronization(nodes, success)) + if (job.completedSynchronization()) { RepairJob completedJob = syncingJobs.remove(job.desc.columnFamily); String remaining = syncingJobs.size() == 0 ? "" : String.format(" (%d remaining column family to sync for this session)", syncingJobs.size()); From f08ae394f0a3ab31260eb0a808160663a857f796 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Fri, 28 Feb 2014 15:11:50 -0600 Subject: [PATCH 2/7] Add CMSClassUnloadingEnabled JVM option Patch by Jonathan Lacefield, reviewed by brandonwilliams for CASSANDRA-6541 --- CHANGES.txt | 1 + conf/cassandra-env.sh | 3 +++ 2 files changed, 4 insertions(+) diff --git a/CHANGES.txt b/CHANGES.txt index 53da840264..780b528773 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 1.2.16 + * Add CMSClassUnloadingEnabled JVM option (CASSANDRA-6541) * Catch memtable flush exceptions during shutdown (CASSANDRA-6735) * Don't attempt cross-dc forwarding in mixed-version cluster with 1.1 (CASSANDRA-6732) diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh index af996efdd4..aa4c3dd93a 100644 --- a/conf/cassandra-env.sh +++ b/conf/cassandra-env.sh @@ -160,6 +160,9 @@ then JVM_OPTS="$JVM_OPTS -javaagent:$CASSANDRA_HOME/lib/jamm-0.2.5.jar" fi +# some JVMs will fill up their heap when accessed via JMX, see CASSANDRA-6541 +JVM_OPTS="$JVM_OPTS -XX:+CMSClassUnloadingEnabled" + # enable thread priorities, primarily so we can give periodic tasks # a lower priority to avoid interfering with client workload JVM_OPTS="$JVM_OPTS -XX:+UseThreadPriorities" From b3a9a443433a271fee33bede60d4892e0c8ffb03 Mon Sep 17 00:00:00 2001 From: Yuki Morishita Date: Fri, 28 Feb 2014 15:28:25 -0600 Subject: [PATCH 3/7] Fix NPE on BulkLoader caused by losing StreamEvent patch by yukim; reviewed by sankalp kohli for CASSANDRA-6636 --- CHANGES.txt | 1 + .../apache/cassandra/io/sstable/SSTableLoader.java | 7 +++---- .../org/apache/cassandra/streaming/StreamPlan.java | 11 ++++++++++- .../cassandra/streaming/StreamResultFuture.java | 7 ++++++- src/java/org/apache/cassandra/tools/BulkLoader.java | 7 ++++--- 5 files changed, 24 insertions(+), 9 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index d6c9ae694e..3e73f91067 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -27,6 +27,7 @@ * Optimize single partition batch statements (CASSANDRA-6737) * Disallow post-query re-ordering when paging (CASSANDRA-6722) * Fix potential paging bug with deleted columns (CASSANDRA-6748) + * Fix NPE on BulkLoader caused by losing StreamEvent (CASSANDRA-6636) Merged from 1.2: * Add CMSClassUnloadingEnabled JVM option (CASSANDRA-6541) * Catch memtable flush exceptions during shutdown (CASSANDRA-6735) diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java b/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java index f867317cd8..1ea4c55b5e 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableLoader.java @@ -144,7 +144,7 @@ public class SSTableLoader implements StreamEventHandler return stream(Collections.emptySet()); } - public StreamResultFuture stream(Set toIgnore) + public StreamResultFuture stream(Set toIgnore, StreamEventHandler... listeners) { client.init(keyspace); outputHandler.output("Established connection to initial hosts"); @@ -175,9 +175,8 @@ public class SSTableLoader implements StreamEventHandler plan.transferFiles(remote, streamingDetails.get(remote)); } - StreamResultFuture bulkResult = plan.execute(); - bulkResult.addEventListener(this); - return bulkResult; + plan.listeners(this, listeners); + return plan.execute(); } public void onSuccess(StreamState finalState) {} diff --git a/src/java/org/apache/cassandra/streaming/StreamPlan.java b/src/java/org/apache/cassandra/streaming/StreamPlan.java index 288929c0dc..740ad66459 100644 --- a/src/java/org/apache/cassandra/streaming/StreamPlan.java +++ b/src/java/org/apache/cassandra/streaming/StreamPlan.java @@ -33,6 +33,7 @@ public class StreamPlan { private final UUID planId = UUIDGen.getTimeUUID(); private final String description; + private final List handlers = new ArrayList<>(); // sessions per InetAddress of the other end. private final Map sessions = new HashMap<>(); @@ -121,6 +122,14 @@ public class StreamPlan return this; } + public StreamPlan listeners(StreamEventHandler handler, StreamEventHandler... handlers) + { + this.handlers.add(handler); + if (handlers != null) + Collections.addAll(this.handlers, handlers); + return this; + } + /** * @return true if this plan has no plan to execute */ @@ -136,7 +145,7 @@ public class StreamPlan */ public StreamResultFuture execute() { - return StreamResultFuture.init(planId, description, sessions.values()); + return StreamResultFuture.init(planId, description, sessions.values(), handlers); } /** diff --git a/src/java/org/apache/cassandra/streaming/StreamResultFuture.java b/src/java/org/apache/cassandra/streaming/StreamResultFuture.java index ccd3c92aee..dcffaff54c 100644 --- a/src/java/org/apache/cassandra/streaming/StreamResultFuture.java +++ b/src/java/org/apache/cassandra/streaming/StreamResultFuture.java @@ -75,9 +75,14 @@ public final class StreamResultFuture extends AbstractFuture set(getCurrentState()); } - static StreamResultFuture init(UUID planId, String description, Collection sessions) + static StreamResultFuture init(UUID planId, String description, Collection sessions, Collection listeners) { StreamResultFuture future = createAndRegister(planId, description, sessions); + if (listeners != null) + { + for (StreamEventHandler listener : listeners) + future.addEventListener(listener); + } logger.info("[Stream #{}] Executing streaming plan for {}", planId, description); // start sessions diff --git a/src/java/org/apache/cassandra/tools/BulkLoader.java b/src/java/org/apache/cassandra/tools/BulkLoader.java index 6c157e29a7..37ec635d98 100644 --- a/src/java/org/apache/cassandra/tools/BulkLoader.java +++ b/src/java/org/apache/cassandra/tools/BulkLoader.java @@ -79,7 +79,10 @@ public class BulkLoader StreamResultFuture future = null; try { - future = loader.stream(options.ignores); + if (options.noProgress) + future = loader.stream(options.ignores); + else + future = loader.stream(options.ignores, new ProgressIndicator()); } catch (Exception e) { @@ -94,8 +97,6 @@ public class BulkLoader } handler.output(String.format("Streaming session ID: %s", future.planId)); - if (!options.noProgress) - future.addEventListener(new ProgressIndicator()); try { From 41d8a5f4861c37ff8e22344418e9277236672c1f Mon Sep 17 00:00:00 2001 From: Marcus Eriksson Date: Mon, 3 Mar 2014 11:35:03 +0100 Subject: [PATCH 4/7] Fix resetAndTruncate:ing CompressionMetadata Patch by kvaster, reviewed by marcuse for CASSANDRA-6791 --- CHANGES.txt | 1 + .../compress/CompressedSequentialWriter.java | 2 +- .../CompressedRandomAccessReaderTest.java | 42 +++++++++++++++++++ 3 files changed, 44 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index 3e73f91067..6de11c50bf 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -28,6 +28,7 @@ * Disallow post-query re-ordering when paging (CASSANDRA-6722) * Fix potential paging bug with deleted columns (CASSANDRA-6748) * Fix NPE on BulkLoader caused by losing StreamEvent (CASSANDRA-6636) + * Fix truncating compression metadata (CASSANDRA-6791) Merged from 1.2: * Add CMSClassUnloadingEnabled JVM option (CASSANDRA-6541) * Catch memtable flush exceptions during shutdown (CASSANDRA-6735) diff --git a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java index 54b990fdd9..eef5b170fb 100644 --- a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java +++ b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java @@ -231,7 +231,7 @@ public class CompressedSequentialWriter extends SequentialWriter // truncate data and index file truncate(chunkOffset); - metadataWriter.resetAndTruncate(realMark.nextChunkIndex); + metadataWriter.resetAndTruncate(realMark.nextChunkIndex - 1); } /** diff --git a/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java b/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java index ee32a0e4d5..3c9dfe526a 100644 --- a/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java +++ b/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java @@ -19,11 +19,13 @@ package org.apache.cassandra.io.compress; import java.io.*; +import java.util.Collections; import java.util.Random; import org.junit.Test; import org.apache.cassandra.db.marshal.BytesType; +import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.io.sstable.CorruptSSTableException; import org.apache.cassandra.io.sstable.SSTableMetadata; import org.apache.cassandra.io.util.*; @@ -48,6 +50,46 @@ public class CompressedRandomAccessReaderTest testResetAndTruncate(File.createTempFile("compressed", "1"), true, 10); testResetAndTruncate(File.createTempFile("compressed", "2"), true, CompressionParameters.DEFAULT_CHUNK_LENGTH); } + @Test + public void test6791() throws IOException, ConfigurationException + { + File f = File.createTempFile("compressed6791_", "3"); + String filename = f.getAbsolutePath(); + try + { + + SSTableMetadata.Collector sstableMetadataCollector = SSTableMetadata.createCollector(BytesType.instance).replayPosition(null); + CompressedSequentialWriter writer = new CompressedSequentialWriter(f, filename + ".metadata", false, new CompressionParameters(SnappyCompressor.instance, 32, Collections.emptyMap()), sstableMetadataCollector); + + for (int i = 0; i < 20; i++) + writer.write("x".getBytes()); + + FileMark mark = writer.mark(); + // write enough garbage to create new chunks: + for (int i = 0; i < 40; ++i) + writer.write("y".getBytes()); + + writer.resetAndTruncate(mark); + + for (int i = 0; i < 20; i++) + writer.write("x".getBytes()); + writer.close(); + + CompressedRandomAccessReader reader = CompressedRandomAccessReader.open(filename, new CompressionMetadata(filename + ".metadata", f.length(), true)); + String res = reader.readLine(); + assertEquals(res, "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); + assertEquals(40, res.length()); + } + finally + { + // cleanup + if (f.exists()) + f.delete(); + File metadata = new File(filename+ ".metadata"); + if (metadata.exists()) + metadata.delete(); + } + } private void testResetAndTruncate(File f, boolean compressed, int junkSize) throws IOException { From 64098f7d6f0b122448693a3c6da16af54c99013b Mon Sep 17 00:00:00 2001 From: Marcus Eriksson Date: Mon, 3 Mar 2014 15:03:34 +0100 Subject: [PATCH 5/7] Fix resetAndTruncate:ing CompressionMetadata Patch by kvaster, reviewed by marcuse for CASSANDRA-6791 --- CHANGES.txt | 2 +- .../compress/CompressedSequentialWriter.java | 2 +- .../CompressedRandomAccessReaderTest.java | 43 +++++++++++++++++++ 3 files changed, 45 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 780b528773..b3c0a3573b 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -19,7 +19,7 @@ * Support negative timestamps for CQL3 dates in query string (CASSANDRA-6718) * Avoid NPEs when receiving table changes for an unknown keyspace (CASSANDRA-5631) * Fix bootstrapping when there is no schema (CASSANDRA-6685) - + * Fix truncating compression metadata (CASSANDRA-6791) 1.2.15 * Move handling of migration event source to solve bootstrap race (CASSANDRA-6648) diff --git a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java index 00eb5a7df5..da55e830f5 100644 --- a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java +++ b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java @@ -231,7 +231,7 @@ public class CompressedSequentialWriter extends SequentialWriter // truncate data and index file truncate(chunkOffset); - metadataWriter.resetAndTruncate(realMark.nextChunkIndex); + metadataWriter.resetAndTruncate(realMark.nextChunkIndex - 1); } /** diff --git a/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java b/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java index 830c3e1b7d..678a650138 100644 --- a/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java +++ b/test/unit/org/apache/cassandra/io/compress/CompressedRandomAccessReaderTest.java @@ -19,10 +19,12 @@ package org.apache.cassandra.io.compress; import java.io.*; +import java.util.Collections; import java.util.Random; import org.junit.Test; +import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.io.sstable.CorruptSSTableException; import org.apache.cassandra.io.sstable.SSTableMetadata; import org.apache.cassandra.io.util.*; @@ -94,6 +96,47 @@ public class CompressedRandomAccessReaderTest } } + @Test + public void test6791() throws IOException, ConfigurationException + { + File f = File.createTempFile("compressed6791_", "3"); + String filename = f.getAbsolutePath(); + try + { + + SSTableMetadata.Collector sstableMetadataCollector = SSTableMetadata.createCollector().replayPosition(null); + CompressedSequentialWriter writer = new CompressedSequentialWriter(f, filename + ".metadata", false, new CompressionParameters(SnappyCompressor.instance, 32, Collections.emptyMap()), sstableMetadataCollector); + + for (int i = 0; i < 20; i++) + writer.write("x".getBytes()); + + FileMark mark = writer.mark(); + // write enough garbage to create new chunks: + for (int i = 0; i < 40; i++) + writer.write("y".getBytes()); + + writer.resetAndTruncate(mark); + + for (int i = 0; i < 20; i++) + writer.write("x".getBytes()); + writer.close(); + + CompressedRandomAccessReader reader = CompressedRandomAccessReader.open(filename, new CompressionMetadata(filename + ".metadata", f.length()), false); + String res = reader.readLine(); + assertEquals(res, "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); + assertEquals(40, res.length()); + } + finally + { + // cleanup + if (f.exists()) + f.delete(); + File metadata = new File(filename + ".metadata"); + if (metadata.exists()) + metadata.delete(); + } + } + @Test public void testDataCorruptionDetection() throws IOException { From 6137b381d2d619ff094395ea32cc4730b9969682 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Tue, 4 Mar 2014 11:16:37 +0100 Subject: [PATCH 6/7] Fix UPDATE updating PRIMARY KEY columns implicitely patch by slebresne; reviewed by iamaleksey for CASSANDRA-6782 --- CHANGES.txt | 1 + .../cassandra/cql3/statements/UpdateStatement.java | 12 +++++++----- 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 6de11c50bf..41c89e13f5 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -29,6 +29,7 @@ * Fix potential paging bug with deleted columns (CASSANDRA-6748) * Fix NPE on BulkLoader caused by losing StreamEvent (CASSANDRA-6636) * Fix truncating compression metadata (CASSANDRA-6791) + * Fix UPDATE updating PRIMARY KEY columns implicitly (CASSANDRA-6782) Merged from 1.2: * Add CMSClassUnloadingEnabled JVM option (CASSANDRA-6541) * Catch memtable flush exceptions during shutdown (CASSANDRA-6735) diff --git a/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java b/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java index dcf22ef9de..fc9bb664e6 100644 --- a/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/UpdateStatement.java @@ -51,17 +51,19 @@ public class UpdateStatement extends ModificationStatement CFDefinition cfDef = cfm.getCfDef(); // Inserting the CQL row marker (see #4361) - // We always need to insert a marker, because of the following situation: + // We always need to insert a marker for INSERT, because of the following situation: // CREATE TABLE t ( k int PRIMARY KEY, c text ); // INSERT INTO t(k, c) VALUES (1, 1) // DELETE c FROM t WHERE k = 1; // SELECT * FROM t; - // The last query should return one row (but with c == null). Adding - // the marker with the insert make sure the semantic is correct (while making sure a - // 'DELETE FROM t WHERE k = 1' does remove the row entirely) + // The last query should return one row (but with c == null). Adding the marker with the insert make sure + // the semantic is correct (while making sure a 'DELETE FROM t WHERE k = 1' does remove the row entirely) + // + // We do not insert the marker for UPDATE however, as this amount to updating the columns in the WHERE + // clause which is inintuitive (#6782) // // We never insert markers for Super CF as this would confuse the thrift side. - if (cfDef.isComposite && !cfDef.isCompact && !cfm.isSuper()) + if (type == StatementType.INSERT && cfDef.isComposite && !cfDef.isCompact && !cfm.isSuper()) { ByteBuffer name = builder.copy().add(ByteBufferUtil.EMPTY_BYTE_BUFFER).build(); cf.addColumn(params.makeColumn(name, ByteBufferUtil.EMPTY_BYTE_BUFFER)); From 0ff7f9969a3f3315eb2fb7e115d72d37dda05157 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Tue, 4 Mar 2014 11:30:54 +0100 Subject: [PATCH 7/7] Fix IllegalArgumentException when upgrading from 1.2 with SC patch by slebresne; reviewed by vijay2win for CASSANDRA-6733 --- CHANGES.txt | 2 + .../db/columniterator/IndexedSliceReader.java | 44 ++++++++++++++++--- .../columniterator/SSTableNamesIterator.java | 6 ++- 3 files changed, 45 insertions(+), 7 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 41c89e13f5..8eb10cdd02 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -30,6 +30,8 @@ * Fix NPE on BulkLoader caused by losing StreamEvent (CASSANDRA-6636) * Fix truncating compression metadata (CASSANDRA-6791) * Fix UPDATE updating PRIMARY KEY columns implicitly (CASSANDRA-6782) + * Fix IllegalArgumentException when updating from 1.2 with SuperColumns + (CASSANDRA-6733) Merged from 1.2: * Add CMSClassUnloadingEnabled JVM option (CASSANDRA-6541) * Catch memtable flush exceptions during shutdown (CASSANDRA-6735) diff --git a/src/java/org/apache/cassandra/db/columniterator/IndexedSliceReader.java b/src/java/org/apache/cassandra/db/columniterator/IndexedSliceReader.java index 036d0cf300..b6aa085e3a 100644 --- a/src/java/org/apache/cassandra/db/columniterator/IndexedSliceReader.java +++ b/src/java/org/apache/cassandra/db/columniterator/IndexedSliceReader.java @@ -179,6 +179,34 @@ class IndexedSliceReader extends AbstractIterator implements OnDiskA } } + static int indexFor(SSTableReader sstable, ByteBuffer name, List indexes, AbstractType comparator, boolean reversed, int startIdx) + { + // If it's a super CF and the sstable is from the old format, then the index will contain old format info, i.e. non composite + // SC names. So we need to 1) use only the SC name part of the comparator and 2) extract only that part from 'name' + if (sstable.metadata.isSuper() && sstable.descriptor.version.hasSuperColumns) + { + AbstractType scComparator = SuperColumns.getComparatorFor(sstable.metadata, false); + ByteBuffer scName = SuperColumns.scName(name); + return IndexHelper.indexFor(scName, indexes, scComparator, reversed, startIdx); + } + return IndexHelper.indexFor(name, indexes, comparator, reversed, startIdx); + } + + static ByteBuffer forIndexComparison(SSTableReader sstable, ByteBuffer name) + { + // See indexFor above. + return sstable.metadata.isSuper() && sstable.descriptor.version.hasSuperColumns + ? SuperColumns.scName(name) + : name; + } + + static AbstractType comparatorForIndex(SSTableReader sstable, AbstractType comparator) + { + return sstable.metadata.isSuper() && sstable.descriptor.version.hasSuperColumns + ? SuperColumns.getComparatorFor(sstable.metadata, false) + : comparator; + } + private abstract class BlockFetcher { protected int currentSliceIdx; @@ -219,16 +247,22 @@ class IndexedSliceReader extends AbstractIterator implements OnDiskA return start.remaining() != 0 && comparator.compare(name, start) < 0; } + protected boolean isIndexEntryBeforeSliceStart(ByteBuffer name) + { + ByteBuffer start = currentStart(); + return start.remaining() != 0 && comparatorForIndex(sstable, comparator).compare(name, forIndexComparison(sstable, start)) < 0; + } + protected boolean isColumnBeforeSliceFinish(OnDiskAtom column) { ByteBuffer finish = currentFinish(); return finish.remaining() == 0 || comparator.compare(column.name(), finish) <= 0; } - protected boolean isAfterSliceFinish(ByteBuffer name) + protected boolean isIndexEntryAfterSliceFinish(ByteBuffer name) { ByteBuffer finish = currentFinish(); - return finish.remaining() != 0 && comparator.compare(name, finish) > 0; + return finish.remaining() != 0 && comparatorForIndex(sstable, comparator).compare(name, forIndexComparison(sstable, finish)) > 0; } } @@ -259,7 +293,7 @@ class IndexedSliceReader extends AbstractIterator implements OnDiskA { while (++currentSliceIdx < slices.length) { - nextIndexIdx = IndexHelper.indexFor(slices[currentSliceIdx].start, indexes, comparator, reversed, nextIndexIdx); + nextIndexIdx = indexFor(sstable, slices[currentSliceIdx].start, indexes, comparator, reversed, nextIndexIdx); if (nextIndexIdx < 0 || nextIndexIdx >= indexes.size()) // no index block for that slice continue; @@ -268,12 +302,12 @@ class IndexedSliceReader extends AbstractIterator implements OnDiskA IndexInfo info = indexes.get(nextIndexIdx); if (reversed) { - if (!isBeforeSliceStart(info.lastName)) + if (!isIndexEntryBeforeSliceStart(info.lastName)) return true; } else { - if (!isAfterSliceFinish(info.firstName)) + if (!isIndexEntryAfterSliceFinish(info.firstName)) return true; } } diff --git a/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java b/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java index 3467244b24..2e84d8d37a 100644 --- a/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java +++ b/src/java/org/apache/cassandra/db/columniterator/SSTableNamesIterator.java @@ -189,13 +189,15 @@ public class SSTableNamesIterator extends AbstractIterator implement int lastIndexIdx = -1; for (ByteBuffer name : columns) { - int index = IndexHelper.indexFor(name, indexList, comparator, false, lastIndexIdx); + int index = IndexedSliceReader.indexFor(sstable, name, indexList, comparator, false, lastIndexIdx); if (index < 0 || index == indexList.size()) continue; IndexHelper.IndexInfo indexInfo = indexList.get(index); // Check the index block does contain the column names and that we haven't inserted this block yet. - if (comparator.compare(name, indexInfo.firstName) < 0 || index == lastIndexIdx) + if (IndexedSliceReader.comparatorForIndex(sstable, comparator).compare(IndexedSliceReader.forIndexComparison(sstable, name), indexInfo.firstName) < 0 + || index == lastIndexIdx) continue; + ranges.add(indexInfo); lastIndexIdx = index; }