diff --git a/src/java/org/apache/cassandra/db/ColumnFamily.java b/src/java/org/apache/cassandra/db/ColumnFamily.java index ddd90c0f1d..2ec0cbe8ce 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamily.java +++ b/src/java/org/apache/cassandra/db/ColumnFamily.java @@ -42,10 +42,8 @@ import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.db.marshal.LongType; -/** - * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) - */ -public final class ColumnFamily + +public final class ColumnFamily implements IColumnContainer { /* The column serializer for this Column Family. Create based on config. */ private static ColumnFamilySerializer serializer_ = new ColumnFamilySerializer(); diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 71ba3099aa..7072cd818a 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -45,6 +45,7 @@ import org.apache.commons.lang.StringUtils; import org.apache.commons.lang.ArrayUtils; import org.apache.commons.collections.IteratorUtils; import org.apache.commons.collections.Predicate; +import org.apache.commons.collections.iterators.ReverseListIterator; import org.cliffc.high_scale_lib.NonBlockingHashMap; @@ -1379,14 +1380,15 @@ public final class ColumnFamilyStore implements ColumnFamilyStoreMBean { QueryFilter nameFilter = new NamesQueryFilter(filter.key, new QueryPath(columnFamily_), filter.path.superColumnName); ColumnFamily cf = getColumnFamily(nameFilter); - if (cf != null) - { - for (IColumn column : cf.getSortedColumns()) - { - filter.filterSuperColumn((SuperColumn)column, gcBefore); - } - } - return removeDeleted(cf, gcBefore); + if (cf == null) + return cf; + + assert cf.getSortedColumns().size() == 1; + SuperColumn sc = (SuperColumn)cf.getSortedColumns().iterator().next(); + SuperColumn scFiltered = filter.filterSuperColumn(sc, gcBefore); + ColumnFamily cfFiltered = cf.cloneMeShallow(); + cfFiltered.addColumn(scFiltered); + return cfFiltered; } // we are querying top-level columns, do a merging fetch with indexes. diff --git a/src/java/org/apache/cassandra/db/IColumnContainer.java b/src/java/org/apache/cassandra/db/IColumnContainer.java new file mode 100644 index 0000000000..e7450a9f6c --- /dev/null +++ b/src/java/org/apache/cassandra/db/IColumnContainer.java @@ -0,0 +1,10 @@ +package org.apache.cassandra.db; + +import org.apache.cassandra.db.marshal.AbstractType; + +public interface IColumnContainer +{ + public void addColumn(IColumn column); + + public AbstractType getComparator(); +} diff --git a/src/java/org/apache/cassandra/db/Memtable.java b/src/java/org/apache/cassandra/db/Memtable.java index b56d87147a..e2218262fb 100644 --- a/src/java/org/apache/cassandra/db/Memtable.java +++ b/src/java/org/apache/cassandra/db/Memtable.java @@ -281,7 +281,7 @@ public class Memtable implements Comparable int index; if (filter.start.length == 0 && !filter.isAscending) { - /* assuming the we scan from the largest column in descending order*/ + /* scan from the largest column in descending order */ index = 0; } else diff --git a/src/java/org/apache/cassandra/db/SuperColumn.java b/src/java/org/apache/cassandra/db/SuperColumn.java index 20dc4653e4..3b45642107 100644 --- a/src/java/org/apache/cassandra/db/SuperColumn.java +++ b/src/java/org/apache/cassandra/db/SuperColumn.java @@ -33,11 +33,8 @@ import org.apache.cassandra.io.ICompactSerializer2; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.marshal.MarshalException; -/** - * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) - */ -public final class SuperColumn implements IColumn +public final class SuperColumn implements IColumn, IColumnContainer { private static Logger logger_ = Logger.getLogger(SuperColumn.class); @@ -47,7 +44,6 @@ public final class SuperColumn implements IColumn } private byte[] name_; - // TODO make subcolumn comparator configurable private ConcurrentSkipListMap columns_; private int localDeletionTime = Integer.MIN_VALUE; private long markedForDeleteAt = Long.MIN_VALUE; diff --git a/src/java/org/apache/cassandra/db/filter/IdentityQueryFilter.java b/src/java/org/apache/cassandra/db/filter/IdentityQueryFilter.java index 46e4f316fe..027812b6cb 100644 --- a/src/java/org/apache/cassandra/db/filter/IdentityQueryFilter.java +++ b/src/java/org/apache/cassandra/db/filter/IdentityQueryFilter.java @@ -14,8 +14,9 @@ public class IdentityQueryFilter extends SliceQueryFilter super(key, path, ArrayUtils.EMPTY_BYTE_ARRAY, ArrayUtils.EMPTY_BYTE_ARRAY, true, Integer.MAX_VALUE); } - public void filterSuperColumn(SuperColumn superColumn, int gcBefore) + public SuperColumn filterSuperColumn(SuperColumn superColumn, int gcBefore) { // no filtering done, deliberately + return superColumn; } } diff --git a/src/java/org/apache/cassandra/db/filter/NamesQueryFilter.java b/src/java/org/apache/cassandra/db/filter/NamesQueryFilter.java index 00237ba2be..a07293b78e 100644 --- a/src/java/org/apache/cassandra/db/filter/NamesQueryFilter.java +++ b/src/java/org/apache/cassandra/db/filter/NamesQueryFilter.java @@ -5,10 +5,7 @@ import java.util.*; import org.apache.cassandra.io.SSTableReader; import org.apache.cassandra.utils.ReducingIterator; -import org.apache.cassandra.db.Memtable; -import org.apache.cassandra.db.ColumnFamily; -import org.apache.cassandra.db.IColumn; -import org.apache.cassandra.db.SuperColumn; +import org.apache.cassandra.db.*; import org.apache.cassandra.db.marshal.AbstractType; public class NamesQueryFilter extends QueryFilter @@ -50,7 +47,7 @@ public class NamesQueryFilter extends QueryFilter return new SSTableNamesIterator(sstable.getFilename(), key, getColumnFamilyName(), columns); } - public void filterSuperColumn(SuperColumn superColumn, int gcBefore) + public SuperColumn filterSuperColumn(SuperColumn superColumn, int gcBefore) { for (IColumn column : superColumn.getSubColumns()) { @@ -59,15 +56,16 @@ public class NamesQueryFilter extends QueryFilter superColumn.remove(column.name()); } } + return superColumn; } - public void collectReducedColumns(ColumnFamily returnCF, Iterator reducedColumns, int gcBefore) + public void collectReducedColumns(IColumnContainer container, Iterator reducedColumns, int gcBefore) { while (reducedColumns.hasNext()) { IColumn column = reducedColumns.next(); if (!column.isMarkedForDelete() || column.getLocalDeletionTime() > gcBefore) - returnCF.addColumn(column); + container.addColumn(column); } } } diff --git a/src/java/org/apache/cassandra/db/filter/QueryFilter.java b/src/java/org/apache/cassandra/db/filter/QueryFilter.java index 984ca8047e..5c97ae18c6 100644 --- a/src/java/org/apache/cassandra/db/filter/QueryFilter.java +++ b/src/java/org/apache/cassandra/db/filter/QueryFilter.java @@ -38,13 +38,13 @@ public abstract class QueryFilter * by the filter code, which should have some limit on the number of columns * to avoid running out of memory on large rows. */ - public abstract void collectReducedColumns(ColumnFamily returnCF, Iterator reducedColumns, int gcBefore); + public abstract void collectReducedColumns(IColumnContainer container, Iterator reducedColumns, int gcBefore); /** * subcolumns of a supercolumn are unindexed, so to pick out parts of those we operate in-memory. - * @param superColumn + * @param superColumn may be modified by filtering op. */ - public abstract void filterSuperColumn(SuperColumn superColumn, int gcBefore); + public abstract SuperColumn filterSuperColumn(SuperColumn superColumn, int gcBefore); public Comparator getColumnComparator(final AbstractType comparator) { diff --git a/src/java/org/apache/cassandra/db/filter/SliceQueryFilter.java b/src/java/org/apache/cassandra/db/filter/SliceQueryFilter.java index cd01ed32f4..c0936ebbaa 100644 --- a/src/java/org/apache/cassandra/db/filter/SliceQueryFilter.java +++ b/src/java/org/apache/cassandra/db/filter/SliceQueryFilter.java @@ -2,13 +2,14 @@ package org.apache.cassandra.db.filter; import java.io.IOException; import java.util.Comparator; -import java.util.Arrays; import java.util.Iterator; +import java.util.ArrayList; +import java.util.List; import org.apache.commons.collections.comparators.ReverseComparator; +import org.apache.commons.collections.iterators.ReverseListIterator; import org.apache.cassandra.io.SSTableReader; -import org.apache.cassandra.utils.ReducingIterator; import org.apache.cassandra.db.*; import org.apache.cassandra.db.marshal.AbstractType; @@ -37,28 +38,21 @@ public class SliceQueryFilter extends QueryFilter return new SSTableSliceIterator(sstable.getFilename(), key, comparator, start, isAscending); } - public void filterSuperColumn(SuperColumn superColumn, int gcBefore) + public SuperColumn filterSuperColumn(SuperColumn superColumn, int gcBefore) { - int liveColumns = 0; - - for (IColumn column : superColumn.getSubColumns()) + SuperColumn scFiltered = superColumn.cloneMeShallow(); + Iterator subcolumns; + if (isAscending) { - final boolean outOfRange = isAscending - ? (start.length > 0 && superColumn.getComparator().compare(column.name(), start) < 0) - || (finish.length > 0 && superColumn.getComparator().compare(column.name(), finish) > 0) - : (start.length > 0 && superColumn.getComparator().compare(column.name(), start) > 0) - || (finish.length > 0 && superColumn.getComparator().compare(column.name(), finish) < 0); - if (outOfRange - || (column.isMarkedForDelete() && column.getLocalDeletionTime() <= gcBefore) - || liveColumns > count) - { - superColumn.remove(column.name()); - } - else if (!column.isMarkedForDelete()) - { - liveColumns++; - } + subcolumns = superColumn.getSubColumns().iterator(); } + else + { + List columnsAsList = new ArrayList(superColumn.getSubColumns()); + subcolumns = new ReverseListIterator(columnsAsList); + } + collectReducedColumns(scFiltered, subcolumns, gcBefore); + return scFiltered; } @Override @@ -67,10 +61,10 @@ public class SliceQueryFilter extends QueryFilter return isAscending ? super.getColumnComparator(comparator) : new ReverseComparator(super.getColumnComparator(comparator)); } - public void collectReducedColumns(ColumnFamily returnCF, Iterator reducedColumns, int gcBefore) + public void collectReducedColumns(IColumnContainer container, Iterator reducedColumns, int gcBefore) { int liveColumns = 0; - AbstractType comparator = returnCF.getComparator(); + AbstractType comparator = container.getComparator(); while (reducedColumns.hasNext()) { @@ -86,7 +80,7 @@ public class SliceQueryFilter extends QueryFilter liveColumns++; if (!column.isMarkedForDelete() || column.getLocalDeletionTime() > gcBefore) - returnCF.addColumn(column); + container.addColumn(column); } } } diff --git a/test/system/test_server.py b/test/system/test_server.py index e12be3e22b..bb602e5c93 100644 --- a/test/system/test_server.py +++ b/test/system/test_server.py @@ -167,6 +167,19 @@ class TestMutations(CassandraTester): _insert_super() _verify_super() + def test_super_subcolumn_limit(self): + _insert_super() + p = SlicePredicate(slice_range=SliceRange('', '', True, 1)) + column_parent = ColumnParent('Super1', 'sc2') + slice = [result.column + for result in client.get_slice('Keyspace1', 'key1', column_parent, p, ConsistencyLevel.ONE)] + assert slice == [Column(_i64(5), 'value5', 0)], slice + p = SlicePredicate(slice_range=SliceRange('', '', False, 1)) + slice = [result.column + for result in client.get_slice('Keyspace1', 'key1', column_parent, p, ConsistencyLevel.ONE)] + assert slice == [Column(_i64(6), 'value6', 0)], slice + + def test_batch_insert(self): _insert_batch(False) time.sleep(0.1)