mirror of https://github.com/apache/cassandra
use the existing collectReducedColumns api to make subcolumn slices conform as expected to filter semantics
patch by jbellis; reviewed by Evan Weaver for CASSANDRA-356 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@802827 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
parent
cee37599f6
commit
ceb4a10808
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
}
|
||||
|
|
@ -281,7 +281,7 @@ public class Memtable implements Comparable<Memtable>
|
|||
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
|
||||
|
|
|
|||
|
|
@ -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<byte[], IColumn> columns_;
|
||||
private int localDeletionTime = Integer.MIN_VALUE;
|
||||
private long markedForDeleteAt = Long.MIN_VALUE;
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<IColumn> reducedColumns, int gcBefore)
|
||||
public void collectReducedColumns(IColumnContainer container, Iterator<IColumn> reducedColumns, int gcBefore)
|
||||
{
|
||||
while (reducedColumns.hasNext())
|
||||
{
|
||||
IColumn column = reducedColumns.next();
|
||||
if (!column.isMarkedForDelete() || column.getLocalDeletionTime() > gcBefore)
|
||||
returnCF.addColumn(column);
|
||||
container.addColumn(column);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<IColumn> reducedColumns, int gcBefore);
|
||||
public abstract void collectReducedColumns(IColumnContainer container, Iterator<IColumn> 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<IColumn> getColumnComparator(final AbstractType comparator)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -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<IColumn> 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<IColumn> columnsAsList = new ArrayList<IColumn>(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<IColumn> reducedColumns, int gcBefore)
|
||||
public void collectReducedColumns(IColumnContainer container, Iterator<IColumn> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Reference in New Issue