Arena allocation for memtables

patch by jbellis and stuhood for CASSANDRA-2252

git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@1158927 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
Jonathan Ellis 2011-08-17 21:52:41 +00:00
parent 9d63eb0647
commit b95a49c5f9
46 changed files with 699 additions and 255 deletions

View File

@ -361,7 +361,7 @@ url=${svn.entry.url}?pathrev=${svn.entry.commit.revision}
</dependency>
<dependency groupId="com.googlecode.json-simple" artifactId="json-simple" version="1.1"/>
<dependency groupId="com.github.stephenc.high-scale-lib" artifactId="high-scale-lib" version="1.1.2"/>
<dependency groupId="com.github.stephenc" artifactId="jamm" version="0.2.2"/>
<dependency groupId="com.github.stephenc" artifactId="jamm" version="0.2.4"/>
<dependency groupId="org.yaml" artifactId="snakeyaml" version="1.6"/>
<dependency groupId="org.apache.cassandra.deps" artifactId="avro" version="1.4.0-cassandra-1">
<exclusion groupId="org.jboss.netty" artifactId="netty"/>
@ -948,7 +948,7 @@ url=${svn.entry.url}?pathrev=${svn.entry.commit.revision}
<jvmarg value="-Dstorage-config=${test.conf}"/>
<jvmarg value="-Daccess.properties=${test.conf}/access.properties"/>
<jvmarg value="-Dlog4j.configuration=log4j-junit.properties" />
<jvmarg value="-javaagent:${basedir}/lib/jamm-0.2.2.jar" />
<jvmarg value="-javaagent:${basedir}/lib/jamm-0.2.4.jar" />
<jvmarg value="-ea"/>
<optjvmargs/>
<classpath>

View File

@ -100,7 +100,7 @@ JVM_OPTS="$JVM_OPTS -ea"
check_openjdk=`"${JAVA:-java}" -version 2>&1 | awk '{if (NR == 2) {print $1}}'`
if [ "$check_openjdk" != "OpenJDK" ]
then
JVM_OPTS="$JVM_OPTS -javaagent:$CASSANDRA_HOME/lib/jamm-0.2.2.jar"
JVM_OPTS="$JVM_OPTS -javaagent:$CASSANDRA_HOME/lib/jamm-0.2.4.jar"
fi
# enable thread priorities, primarily so we can give periodic tasks

Binary file not shown.

BIN
lib/jamm-0.2.4.jar Normal file

Binary file not shown.

View File

@ -35,7 +35,9 @@ import org.apache.cassandra.db.filter.QueryPath;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.io.ICompactSerializer2;
import org.apache.cassandra.io.util.IIterableColumns;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.HeapAllocator;
public abstract class AbstractColumnContainer implements IColumnContainer, IIterableColumns
{
@ -105,15 +107,20 @@ public abstract class AbstractColumnContainer implements IColumnContainer, IIter
/**
* We need to go through each column in the column container and resolve it before adding
*/
public void addAll(AbstractColumnContainer cc)
public void addAll(AbstractColumnContainer cc, Allocator allocator)
{
columns.addAll(cc.columns);
columns.addAll(cc.columns, allocator);
delete(cc);
}
public void addColumn(IColumn column)
{
columns.addColumn(column);
addColumn(column, HeapAllocator.instance);
}
public void addColumn(IColumn column, Allocator allocator)
{
columns.addColumn(column, allocator);
}
public IColumn getColumn(ByteBuffer name)

View File

@ -21,6 +21,7 @@ import java.nio.ByteBuffer;
import java.util.*;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.utils.Allocator;
/**
* A ISortedColumns backed by an ArrayList.
@ -100,7 +101,7 @@ public class ArrayBackedSortedColumns extends ArrayList<IColumn> implements ISor
* without knowing about (we can revisit that decision later if we have
* use cases where most insert are in sorted order but a few are not).
*/
public void addColumn(IColumn column)
public void addColumn(IColumn column, Allocator allocator)
{
if (isEmpty())
{
@ -122,13 +123,13 @@ public class ArrayBackedSortedColumns extends ArrayList<IColumn> implements ISor
else if (c == 0)
{
// Resolve against last
resolveAgainst(size() - 1, column);
resolveAgainst(size() - 1, column, allocator);
}
else
{
int pos = binarySearch(column.name());
if (pos >= 0)
resolveAgainst(pos, column);
resolveAgainst(pos, column, allocator);
else
add(-pos-1, column);
}
@ -138,19 +139,19 @@ public class ArrayBackedSortedColumns extends ArrayList<IColumn> implements ISor
* Resolve against element at position i.
* Assume that i is a valid position.
*/
private void resolveAgainst(int i, IColumn column)
private void resolveAgainst(int i, IColumn column, Allocator allocator)
{
IColumn oldColumn = get(i);
if (oldColumn instanceof SuperColumn)
{
// Delegated to SuperColumn
assert column instanceof SuperColumn;
((SuperColumn) oldColumn).putColumn((SuperColumn)column);
((SuperColumn) oldColumn).putColumn((SuperColumn)column, allocator);
}
else
{
// calculate reconciled col from old (existing) col and new col
IColumn reconciledColumn = column.reconcile(oldColumn);
IColumn reconciledColumn = column.reconcile(oldColumn, allocator);
set(i, reconciledColumn);
}
}
@ -186,7 +187,7 @@ public class ArrayBackedSortedColumns extends ArrayList<IColumn> implements ISor
return -mid - (result < 0 ? 1 : 2);
}
public void addAll(ISortedColumns cm)
public void addAll(ISortedColumns cm, Allocator allocator)
{
if (cm.isEmpty())
return;
@ -214,7 +215,7 @@ public class ArrayBackedSortedColumns extends ArrayList<IColumn> implements ISor
else // c == 0
{
add(copy[idx]);
resolveAgainst(size() - 1, otherColumn);
resolveAgainst(size() - 1, otherColumn, allocator);
idx++;
otherColumn = other.hasNext() ? other.next() : null;
}

View File

@ -30,8 +30,9 @@ import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
/**
* Column is immutable, which prevents all kinds of confusion in a multithreaded environment.
@ -146,6 +147,11 @@ public class Column implements IColumn
}
public void addColumn(IColumn column)
{
addColumn(null, null);
}
public void addColumn(IColumn column, Allocator allocator)
{
throw new UnsupportedOperationException("This operation is not supported for simple columns.");
}
@ -183,6 +189,11 @@ public class Column implements IColumn
}
public IColumn reconcile(IColumn column)
{
return reconcile(column, HeapAllocator.instance);
}
public IColumn reconcile(IColumn column, Allocator allocator)
{
// tombstones take precedence. (if both are tombstones, then it doesn't matter which one we use.)
if (isMarkedForDelete())
@ -225,9 +236,14 @@ public class Column implements IColumn
public IColumn localCopy(ColumnFamilyStore cfs)
{
return new Column(cfs.internOrCopy(name), ByteBufferUtil.clone(value), timestamp);
return localCopy(cfs, HeapAllocator.instance);
}
public IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator)
{
return new Column(cfs.internOrCopy(name, allocator), allocator.clone(value), timestamp);
}
public String getString(AbstractType comparator)
{
StringBuilder sb = new StringBuilder();

View File

@ -31,7 +31,9 @@ import org.apache.cassandra.db.filter.QueryPath;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.io.IColumnSerializer;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.HeapAllocator;
public class ColumnFamily extends AbstractColumnContainer
{
@ -121,6 +123,10 @@ public class ColumnFamily extends AbstractColumnContainer
return cfm;
}
/**
* FIXME: shouldn't need to hold a reference to a serializer; worse, for super cfs,
* it will be a _unique_ serializer object per row
*/
public IColumnSerializer getColumnSerializer()
{
return cfm.getColumnSerializer();
@ -301,11 +307,16 @@ public class ColumnFamily extends AbstractColumnContainer
}
public void resolve(ColumnFamily cf)
{
resolve(cf, HeapAllocator.instance);
}
public void resolve(ColumnFamily cf, Allocator allocator)
{
// Row _does_ allow null CF objects :( seems a necessary evil for efficiency
if (cf == null)
return;
addAll(cf);
addAll(cf, allocator);
}
public long serializedSize()

View File

@ -110,10 +110,10 @@ public class ColumnFamilySerializer implements ICompactSerializer3<ColumnFamily>
public ColumnFamily deserialize(DataInput dis) throws IOException
{
return deserialize(dis, false, false, ThreadSafeSortedColumns.FACTORY);
return deserialize(dis, false, ThreadSafeSortedColumns.FACTORY);
}
public ColumnFamily deserialize(DataInput dis, boolean intern, boolean fromRemote, ISortedColumns.Factory factory) throws IOException
public ColumnFamily deserialize(DataInput dis, boolean fromRemote, ISortedColumns.Factory factory) throws IOException
{
if (!dis.readBoolean())
return null;
@ -124,23 +124,22 @@ public class ColumnFamilySerializer implements ICompactSerializer3<ColumnFamily>
throw new UnserializableColumnFamilyException("Couldn't find cfId=" + cfId, cfId);
ColumnFamily cf = ColumnFamily.create(cfId, factory);
deserializeFromSSTableNoColumns(cf, dis);
deserializeColumns(dis, cf, intern, fromRemote);
deserializeColumns(dis, cf, fromRemote);
return cf;
}
public void deserializeColumns(DataInput dis, ColumnFamily cf, boolean intern, boolean fromRemote) throws IOException
public void deserializeColumns(DataInput dis, ColumnFamily cf, boolean fromRemote) throws IOException
{
int count = dis.readInt();
deserializeColumns(dis, cf, count, intern, fromRemote);
int size = dis.readInt();
deserializeColumns(dis, cf, size, fromRemote);
}
/* column count is already read from DataInput */
public void deserializeColumns(DataInput dis, ColumnFamily cf, int count, boolean intern, boolean fromRemote) throws IOException
public void deserializeColumns(DataInput dis, ColumnFamily cf, int size, boolean fromRemote) throws IOException
{
ColumnFamilyStore interner = intern ? Table.open(CFMetaData.getCF(cf.id()).left).getColumnFamilyStore(cf.id()) : null;
for (int i = 0; i < count; ++i)
for (int i = 0; i < size; ++i)
{
IColumn column = cf.getColumnSerializer().deserialize(dis, interner, fromRemote, (int) (System.currentTimeMillis() / 1000));
IColumn column = cf.getColumnSerializer().deserialize(dis, fromRemote, (int) (System.currentTimeMillis() / 1000));
cf.addColumn(column);
}
}

View File

@ -749,7 +749,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
{
ColumnFamily cachedRow = getRawCachedRow(key);
if (cachedRow != null)
cachedRow.addAll(columnFamily);
cachedRow.addAll(columnFamily, HeapAllocator.instance);
}
}
@ -1211,7 +1211,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
{
ColumnFamily cf = cached.cloneMeShallow();
if (sc != null)
cf.addColumn(sc);
cf.addColumn(sc, HeapAllocator.instance);
return removeDeleted(cf, gcBefore);
}
}
@ -1919,10 +1919,10 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
return internedName;
}
public ByteBuffer internOrCopy(ByteBuffer name)
public ByteBuffer internOrCopy(ByteBuffer name, Allocator allocator)
{
if (internedNames.size() >= INTERN_CUTOFF)
return ByteBufferUtil.clone(name);
return allocator.clone(name);
return intern(name);
}
@ -1930,7 +1930,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
public ByteBuffer maybeIntern(ByteBuffer name)
{
if (internedNames.size() >= INTERN_CUTOFF)
return name;
return null;
return intern(name);
}

View File

@ -69,7 +69,7 @@ public class ColumnSerializer implements IColumnSerializer
public Column deserialize(DataInput dis) throws IOException
{
return deserialize(dis, null, false);
return deserialize(dis, false);
}
/*
@ -77,18 +77,16 @@ public class ColumnSerializer implements IColumnSerializer
* deserialize comes from a remote host. If it does, then we must clear
* the delta.
*/
public Column deserialize(DataInput dis, ColumnFamilyStore interner, boolean fromRemote) throws IOException
public Column deserialize(DataInput dis, boolean fromRemote) throws IOException
{
return deserialize(dis, interner, fromRemote, (int) (System.currentTimeMillis() / 1000));
return deserialize(dis, fromRemote, (int) (System.currentTimeMillis() / 1000));
}
public Column deserialize(DataInput dis, ColumnFamilyStore interner, boolean fromRemote, int expireBefore) throws IOException
public Column deserialize(DataInput dis, boolean fromRemote, int expireBefore) throws IOException
{
ByteBuffer name = ByteBufferUtil.readWithShortLength(dis);
if (name.remaining() <= 0)
throw new CorruptColumnException("invalid column name length " + name.remaining());
if (interner != null)
name = interner.maybeIntern(name);
int b = dis.readUnsignedByte();
if ((b & COUNTER_MASK) != 0)

View File

@ -31,7 +31,9 @@ import org.apache.cassandra.db.context.IContext.ContextRelationship;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
import org.apache.cassandra.utils.NodeId;
/**
@ -47,12 +49,12 @@ public class CounterColumn extends Column
public CounterColumn(ByteBuffer name, long value, long timestamp)
{
this(name, contextManager.create(value), timestamp);
this(name, contextManager.create(value, HeapAllocator.instance), timestamp);
}
public CounterColumn(ByteBuffer name, long value, long timestamp, long timestampOfLastDelete)
{
this(name, contextManager.create(value), timestamp, timestampOfLastDelete);
this(name, contextManager.create(value, HeapAllocator.instance), timestamp, timestampOfLastDelete);
}
public CounterColumn(ByteBuffer name, ByteBuffer value, long timestamp)
@ -135,7 +137,7 @@ public class CounterColumn extends Column
}
@Override
public IColumn reconcile(IColumn column)
public IColumn reconcile(IColumn column, Allocator allocator)
{
assert (column instanceof CounterColumn) || (column instanceof DeletedColumn) : "Wrong class type.";
@ -162,7 +164,7 @@ public class CounterColumn extends Column
// live + live: merge clocks; update value
return new CounterColumn(
name(),
contextManager.merge(value(), column.value()),
contextManager.merge(value(), column.value(), allocator),
Math.max(timestamp(), column.timestamp()),
Math.max(timestampOfLastDelete(), ((CounterColumn)column).timestampOfLastDelete()));
}
@ -185,7 +187,13 @@ public class CounterColumn extends Column
@Override
public IColumn localCopy(ColumnFamilyStore cfs)
{
return new CounterColumn(cfs.internOrCopy(name), ByteBufferUtil.clone(value), timestamp, timestampOfLastDelete);
return new CounterColumn(cfs.internOrCopy(name, HeapAllocator.instance), ByteBufferUtil.clone(value), timestamp, timestampOfLastDelete);
}
@Override
public IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator)
{
return new CounterColumn(cfs.internOrCopy(name, allocator), allocator.clone(value), timestamp, timestampOfLastDelete);
}
@Override

View File

@ -40,7 +40,9 @@ import org.apache.cassandra.io.ICompactSerializer;
import org.apache.cassandra.io.util.FastByteArrayOutputStream;
import org.apache.cassandra.net.Message;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.thrift.ConsistencyLevel;
@ -144,7 +146,7 @@ public class CounterMutation implements IMutation
localMutation.add(merger);
localMutation.apply();
cf.addAll(merger);
cf.addAll(merger, HeapAllocator.instance);
}
}
return row;
@ -216,8 +218,7 @@ public class CounterMutation implements IMutation
public void apply() throws IOException
{
// We need to transform all CounterUpdateColumn to CounterColumn and we need to deepCopy. Both are done
// below since CUC.asCounterColumn() does a deep copy.
// transform all CounterUpdateColumn to CounterColumn: accomplished by localCopy
RowMutation rm = new RowMutation(rowMutation.getTable(), ByteBufferUtil.clone(rowMutation.key()));
Table table = Table.open(rm.getTable());
@ -227,7 +228,7 @@ public class CounterMutation implements IMutation
ColumnFamilyStore cfs = table.getColumnFamilyStore(cf.id());
for (IColumn column : cf_)
{
cf.addColumn(column.localCopy(cfs));
cf.addColumn(column.localCopy(cfs), HeapAllocator.instance);
}
rm.add(cf);
}

View File

@ -21,7 +21,9 @@ package org.apache.cassandra.db;
import java.nio.ByteBuffer;
import org.apache.cassandra.db.context.CounterContext;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
/**
* A counter update while it hasn't been applied yet by the leader replica.
@ -55,7 +57,7 @@ public class CounterUpdateColumn extends Column
}
@Override
public IColumn reconcile(IColumn column)
public IColumn reconcile(IColumn column, Allocator allocator)
{
// The only time this could happen is if a batchAdd ships two
// increment for the same column. Hence we simply sums the delta.
@ -80,8 +82,17 @@ public class CounterUpdateColumn extends Column
@Override
public CounterColumn localCopy(ColumnFamilyStore cfs)
{
return new CounterColumn(cfs.internOrCopy(name),
CounterContext.instance().create(delta()),
return new CounterColumn(cfs.internOrCopy(name, HeapAllocator.instance),
CounterContext.instance().create(delta(), HeapAllocator.instance),
timestamp(),
Long.MIN_VALUE);
}
@Override
public IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator)
{
return new CounterColumn(cfs.internOrCopy(name, allocator),
CounterContext.instance().create(delta(), allocator),
timestamp(),
Long.MIN_VALUE);
}

View File

@ -22,7 +22,9 @@ import java.nio.ByteBuffer;
import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
public class DeletedColumn extends Column
{
@ -55,17 +57,23 @@ public class DeletedColumn extends Column
}
@Override
public IColumn reconcile(IColumn column)
public IColumn reconcile(IColumn column, Allocator allocator)
{
if (column instanceof DeletedColumn)
return super.reconcile(column);
return column.reconcile(this);
return super.reconcile(column, allocator);
return column.reconcile(this, allocator);
}
@Override
public IColumn localCopy(ColumnFamilyStore cfs)
{
return new DeletedColumn(cfs.internOrCopy(name), ByteBufferUtil.clone(value), timestamp);
return new DeletedColumn(cfs.internOrCopy(name, HeapAllocator.instance), ByteBufferUtil.clone(value), timestamp);
}
@Override
public IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator)
{
return new DeletedColumn(cfs.internOrCopy(name, allocator), allocator.clone(value), timestamp);
}
@Override

View File

@ -26,7 +26,9 @@ import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
/**
* Alternative to Column that have an expiring time.
@ -117,7 +119,16 @@ public class ExpiringColumn extends Column
@Override
public IColumn localCopy(ColumnFamilyStore cfs)
{
return new ExpiringColumn(cfs.internOrCopy(name), ByteBufferUtil.clone(value), timestamp, timeToLive, localExpirationTime);
return new ExpiringColumn(cfs.internOrCopy(name, HeapAllocator.instance), ByteBufferUtil.clone(value), timestamp, timeToLive, localExpirationTime);
}
@Override
public IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator)
{
ByteBuffer clonedName = cfs.maybeIntern(name);
if (clonedName == null)
clonedName = allocator.clone(name);
return new ExpiringColumn(clonedName, allocator.clone(value), timestamp, timeToLive, localExpirationTime);
}
@Override

View File

@ -25,8 +25,10 @@ import java.util.Collection;
import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.FBUtilities;
/** TODO: rename */
public interface IColumn
{
public static final int MAX_NAME_LENGTH = FBUtilities.MAX_UNSIGNED_SHORT;
@ -43,17 +45,25 @@ public interface IColumn
public Collection<IColumn> getSubColumns();
public IColumn getSubColumn(ByteBuffer columnName);
public void addColumn(IColumn column);
public void addColumn(IColumn column, Allocator allocator);
public IColumn diff(IColumn column);
public IColumn reconcile(IColumn column);
public IColumn reconcile(IColumn column, Allocator allocator);
public void updateDigest(MessageDigest digest);
public int getLocalDeletionTime(); // for tombstone GC, so int is sufficient granularity
public String getString(AbstractType comparator);
public void validateFields(CFMetaData metadata) throws MarshalException;
/** clones the column, interning column names and making copies of other underlying byte buffers
* @param cfs*/
/** clones the column for the row cache, interning column names and making copies of other underlying byte buffers */
IColumn localCopy(ColumnFamilyStore cfs);
/**
* clones the column for the memtable, interning column names and making copies of other underlying byte buffers.
* Unlike the other localCopy, this uses Allocator to allocate values in contiguous memory regions,
* which helps avoid heap fragmentation.
*/
IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator);
/**
* For a simple column, live == !isMarkedForDelete.
* For a supercolumn, live means it has at least one subcolumn whose timestamp is greater than the

View File

@ -24,10 +24,12 @@ import java.nio.ByteBuffer;
import java.util.Collection;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.utils.Allocator;
public interface IColumnContainer
{
public void addColumn(IColumn column);
public void addColumn(IColumn column, Allocator allocator);
public void remove(ByteBuffer columnName);
/**

View File

@ -25,6 +25,7 @@ import java.util.SortedSet;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.io.util.IIterableColumns;
import org.apache.cassandra.utils.Allocator;
/**
* A sorted map of columns.
@ -45,7 +46,7 @@ public interface ISortedColumns extends IIterableColumns
* If a column with the same name is already present in the map, it will
* be replaced by the newly added column.
*/
public void addColumn(IColumn column);
public void addColumn(IColumn column, Allocator allocator);
/**
* Adds all the columns of a given column map to this column map.
@ -56,7 +57,7 @@ public interface ISortedColumns extends IIterableColumns
* </code>
* but is potentially faster.
*/
public void addAll(ISortedColumns cm);
public void addAll(ISortedColumns cm, Allocator allocator);
/**
* Replace oldColumn if present by newColumn.

View File

@ -43,6 +43,7 @@ import org.apache.cassandra.db.filter.SliceQueryFilter;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.io.sstable.SSTableReader;
import org.apache.cassandra.io.sstable.SSTableWriter;
import org.apache.cassandra.utils.SlabAllocator;
import org.apache.cassandra.utils.WrappedRunnable;
import org.github.jamm.MemoryMeter;
@ -54,7 +55,7 @@ public class Memtable
private static final double MIN_SANE_LIVE_RATIO = 1.0;
// max liveratio seen w/ 1-byte columns on a 64-bit jvm was 19. If it gets higher than 64 something is probably broken.
private static final double MAX_SANE_LIVE_RATIO = 64.0;
private static final MemoryMeter meter = new MemoryMeter();
private static final MemoryMeter meter = new MemoryMeter().omitSharedBufferOverhead();
// we're careful to only allow one count to run at a time because counting is slow
// (can be minutes, for a large memtable and a busy server), so we could keep memtables
@ -69,6 +70,8 @@ public class Memtable
}
};
volatile static Memtable activelyMeasuring;
private volatile boolean isFrozen;
private final AtomicLong currentThroughput = new AtomicLong(0);
private final AtomicLong currentOperations = new AtomicLong(0);
@ -78,7 +81,7 @@ public class Memtable
private final long THRESHOLD;
private final long THRESHOLD_COUNT;
volatile static Memtable activelyMeasuring;
private SlabAllocator allocator = new SlabAllocator();
public Memtable(ColumnFamilyStore cfs)
{
@ -89,8 +92,9 @@ public class Memtable
public long getLiveSize()
{
// 25% fudge factor
return (long) (currentThroughput.get() * cfs.liveRatio * 1.25);
// 25% fudge factor on the base throughput * liveRatio calculation. (Based on observed
// pre-slabbing behavior -- not sure what accounts for this. May have changed with introduction of slabbing.)
return (long) (currentThroughput.get() * cfs.liveRatio * 1.25) + allocator.size();
}
public long getSerializedSize()
@ -192,11 +196,25 @@ public class Memtable
? cf.isMarkedForDelete() ? 1 : 0
: cf.getColumnCount());
ColumnFamily oldCf = columnFamilies.putIfAbsent(key, cf);
if (oldCf == null)
return;
ColumnFamily clonedCf = columnFamilies.get(key);
// if the row doesn't exist yet in the memtable, clone cf to our allocator.
if (clonedCf == null)
{
clonedCf = cf.cloneMeShallow();
for (IColumn column : cf.getSortedColumns())
clonedCf.addColumn(column.localCopy(cfs, allocator));
clonedCf = columnFamilies.putIfAbsent(new DecoratedKey(key.token, allocator.clone(key.key)), clonedCf);
if (clonedCf == null)
return;
// else there was a race and the other thread won. fall through to updating his CF object
}
oldCf.resolve(cf);
// we duplicate the funcationality of CF.resolve here to avoid having to either pass the Memtable in for
// the cloning operation, or cloning the CF container as well as the Columns. fortunately, resolve
// is really quite simple:
clonedCf.delete(cf);
for (IColumn column : cf.getSortedColumns())
clonedCf.addColumn(column.localCopy(cfs, allocator), allocator);
}
// for debugging

View File

@ -66,7 +66,7 @@ public class Row
public Row deserialize(DataInputStream dis, int version, boolean fromRemote, ISortedColumns.Factory factory) throws IOException
{
return new Row(StorageService.getPartitioner().decorateKey(ByteBufferUtil.readWithShortLength(dis)),
ColumnFamily.serializer().deserialize(dis, false, fromRemote, factory));
ColumnFamily.serializer().deserialize(dis, fromRemote, factory));
}
public Row deserialize(DataInputStream dis, int version) throws IOException

View File

@ -373,23 +373,6 @@ public class RowMutation implements IMutation, MessageProducer
return rm;
}
public RowMutation localCopy()
{
RowMutation rm = new RowMutation(table_, ByteBufferUtil.clone(key_));
Table table = Table.open(table_);
for (Map.Entry<Integer, ColumnFamily> entry : modifications_.entrySet())
{
ColumnFamily cf = entry.getValue().cloneMeShallow();
ColumnFamilyStore cfs = table.getColumnFamilyStore(cf.id());
for (IColumn col : entry.getValue())
cf.addColumn(col.localCopy(cfs));
rm.modifications_.put(entry.getKey(), cf);
}
return rm;
}
public static class RowMutationSerializer implements ICompactSerializer<RowMutation>
{
public void serialize(RowMutation rm, DataOutputStream dos, int version) throws IOException
@ -419,7 +402,7 @@ public class RowMutation implements IMutation, MessageProducer
for (int i = 0; i < size; ++i)
{
Integer cfid = Integer.valueOf(dis.readInt());
ColumnFamily cf = ColumnFamily.serializer().deserialize(dis, true, fromRemote, ThreadSafeSortedColumns.FACTORY);
ColumnFamily cf = ColumnFamily.serializer().deserialize(dis, fromRemote, ThreadSafeSortedColumns.FACTORY);
modifications.put(cfid, cf);
}
return new RowMutation(table, key, modifications);

View File

@ -25,9 +25,6 @@ import java.nio.ByteBuffer;
import java.security.MessageDigest;
import java.util.Collection;
import java.util.Comparator;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.ConcurrentSkipListMap;
import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.db.marshal.AbstractType;
@ -35,10 +32,11 @@ import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.io.IColumnSerializer;
import org.apache.cassandra.io.util.ColumnSortedMap;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
import org.cliffc.high_scale_lib.NonBlockingHashMap;
public class SuperColumn extends AbstractColumnContainer implements IColumn
{
private static NonBlockingHashMap<Comparator, SuperColumnSerializer> serializers = new NonBlockingHashMap<Comparator, SuperColumnSerializer>();
@ -160,21 +158,21 @@ public class SuperColumn extends AbstractColumnContainer implements IColumn
}
@Override
public void addColumn(IColumn column)
public void addColumn(IColumn column, Allocator allocator)
{
assert column instanceof Column : "A super column can only contain simple columns";
super.addColumn((Column)column);
super.addColumn(column, allocator);
}
/*
* Go through each sub column if it exists then as it to resolve itself
* if the column does not exist then create it.
*/
void putColumn(SuperColumn column)
void putColumn(SuperColumn column, Allocator allocator)
{
for (IColumn subColumn : column.getSubColumns())
{
addColumn(subColumn);
addColumn(subColumn, allocator);
}
delete(column);
}
@ -255,29 +253,33 @@ public class SuperColumn extends AbstractColumnContainer implements IColumn
return mostRecentLiveChangeAt() > getMarkedForDeleteAt();
}
public IColumn shallowCopy()
{
SuperColumn sc = new SuperColumn(ByteBufferUtil.clone(name()), this.getComparator());
// since deletion info is immutable, aliasing it is fine
sc.deletionInfo.set(deletionInfo.get());
return sc;
}
public IColumn localCopy(ColumnFamilyStore cfs)
{
return localCopy(cfs, HeapAllocator.instance);
}
public IColumn localCopy(ColumnFamilyStore cfs, Allocator allocator)
{
// we don't try to intern supercolumn names, because if we're using Cassandra correctly it's almost
// certainly just going to pollute our interning map with unique, dynamic values
SuperColumn sc = (SuperColumn)shallowCopy();
SuperColumn sc = new SuperColumn(allocator.clone(name), this.getComparator());
// since deletion info is immutable, aliasing it is fine
sc.deletionInfo.set(deletionInfo.get());
for(IColumn c : columns)
{
sc.addColumn(c.localCopy(cfs));
sc.addColumn(c.localCopy(cfs, allocator));
}
return sc;
}
public IColumn reconcile(IColumn c)
{
return reconcile(null, null);
}
public IColumn reconcile(IColumn c, Allocator allocator)
{
throw new UnsupportedOperationException("This operation is unsupported on super columns.");
}
@ -335,20 +337,15 @@ class SuperColumnSerializer implements IColumnSerializer
public IColumn deserialize(DataInput dis) throws IOException
{
return deserialize(dis, null, false);
return deserialize(dis, false);
}
public IColumn deserialize(DataInput dis, ColumnFamilyStore interner) throws IOException
public IColumn deserialize(DataInput dis, boolean fromRemote) throws IOException
{
return deserialize(dis, interner, false);
return deserialize(dis, fromRemote, (int)(System.currentTimeMillis() / 1000));
}
public IColumn deserialize(DataInput dis, ColumnFamilyStore interner, boolean fromRemote) throws IOException
{
return deserialize(dis, interner, fromRemote, (int)(System.currentTimeMillis() / 1000));
}
public IColumn deserialize(DataInput dis, ColumnFamilyStore interner, boolean fromRemote, int expireBefore) throws IOException
public IColumn deserialize(DataInput dis, boolean fromRemote, int expireBefore) throws IOException
{
ByteBuffer name = ByteBufferUtil.readWithShortLength(dis);
int localDeleteTime = dis.readInt();
@ -361,7 +358,7 @@ class SuperColumnSerializer implements IColumnSerializer
/* read the number of columns */
int size = dis.readInt();
ColumnSerializer serializer = Column.serializer();
ColumnSortedMap preSortedMap = new ColumnSortedMap(comparator, serializer, dis, interner, size, fromRemote, expireBefore);
ColumnSortedMap preSortedMap = new ColumnSortedMap(comparator, serializer, dis, size, fromRemote, expireBefore);
SuperColumn superColumn = new SuperColumn(name, ThreadSafeSortedColumns.factory().fromSorted(preSortedMap, false));
if (localDeleteTime != Integer.MIN_VALUE && localDeleteTime <= 0)
{

View File

@ -18,15 +18,14 @@
package org.apache.cassandra.db;
import java.nio.ByteBuffer;
import java.util.Comparator;
import java.util.Collection;
import java.util.Iterator;
import java.util.Map;
import java.util.SortedMap;
import java.util.SortedSet;
import java.util.concurrent.ConcurrentSkipListMap;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.utils.Allocator;
public class ThreadSafeSortedColumns extends ConcurrentSkipListMap<ByteBuffer, IColumn> implements ISortedColumns
{
@ -72,7 +71,7 @@ public class ThreadSafeSortedColumns extends ConcurrentSkipListMap<ByteBuffer, I
* If we find an old column that has the same name
* the ask it to resolve itself else add the new column
*/
public void addColumn(IColumn column)
public void addColumn(IColumn column, Allocator allocator)
{
ByteBuffer name = column.name();
IColumn oldColumn;
@ -81,13 +80,13 @@ public class ThreadSafeSortedColumns extends ConcurrentSkipListMap<ByteBuffer, I
if (oldColumn instanceof SuperColumn)
{
assert column instanceof SuperColumn;
((SuperColumn) oldColumn).putColumn((SuperColumn)column);
((SuperColumn) oldColumn).putColumn((SuperColumn)column, allocator);
break; // Delegated to SuperColumn
}
else
{
// calculate reconciled col from old (existing) col and new col
IColumn reconciledColumn = column.reconcile(oldColumn);
IColumn reconciledColumn = column.reconcile(oldColumn, allocator);
if (replace(name, oldColumn, reconciledColumn))
break;
@ -100,10 +99,10 @@ public class ThreadSafeSortedColumns extends ConcurrentSkipListMap<ByteBuffer, I
/**
* We need to go through each column in the column container and resolve it before adding
*/
public void addAll(ISortedColumns cm)
public void addAll(ISortedColumns cm, Allocator allocator)
{
for (IColumn column : cm.getSortedColumns())
addColumn(column);
addColumn(column, allocator);
}
public boolean replace(IColumn oldColumn, IColumn newColumn)

View File

@ -36,6 +36,7 @@ import org.apache.cassandra.db.CounterColumn;
import org.apache.cassandra.db.DecoratedKey;
import org.apache.cassandra.io.sstable.SSTableIdentityIterator;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.utils.HeapAllocator;
/**
* PrecompactedRow merges its rows in its constructor in memory.
@ -96,7 +97,7 @@ public class PrecompactedRow extends AbstractCompactedRow
}
else
{
cf.addAll(thisCF);
cf.addAll(thisCF, HeapAllocator.instance);
}
}
return cf;

View File

@ -25,6 +25,8 @@ import org.apache.log4j.Logger;
import org.apache.cassandra.db.marshal.MarshalException;
import org.apache.cassandra.db.DBConstants;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.HeapAllocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.NodeId;
@ -85,19 +87,21 @@ public class CounterContext implements IContext
}
/**
* Creates an initial counter context with an initial value for the local node with.
* Creates an initial counter context with an initial value for the local node.
*
*
* @param value the value for this initial update
*
* @param allocator
* @return an empty counter context.
*/
public ByteBuffer create(long value)
public ByteBuffer create(long value, Allocator allocator)
{
ByteBuffer context = ByteBuffer.allocate(HEADER_SIZE_LENGTH + HEADER_ELT_LENGTH + STEP_LENGTH);
ByteBuffer context = allocator.allocate(HEADER_SIZE_LENGTH + HEADER_ELT_LENGTH + STEP_LENGTH);
// The first (and only) elt is a delta
context.putShort(0, (short)1);
context.putShort(HEADER_SIZE_LENGTH, (short)0);
writeElementAtOffset(context, HEADER_SIZE_LENGTH + HEADER_ELT_LENGTH, NodeId.getLocalId(), 1L, value);
context.putShort(context.position(), (short)1);
context.putShort(context.position() + HEADER_SIZE_LENGTH, (short)0);
writeElementAtOffset(context, context.position() + HEADER_SIZE_LENGTH + HEADER_ELT_LENGTH, NodeId.getLocalId(), 1L, value);
return context;
}
@ -105,21 +109,23 @@ public class CounterContext implements IContext
public ByteBuffer create(NodeId id, long clock, long value, boolean isDelta)
{
ByteBuffer context = ByteBuffer.allocate(HEADER_SIZE_LENGTH + (isDelta ? HEADER_ELT_LENGTH : 0) + STEP_LENGTH);
context.putShort(0, (short)(isDelta ? 1 : 0));
context.putShort(context.position(), (short)(isDelta ? 1 : 0));
if (isDelta)
{
context.putShort(HEADER_SIZE_LENGTH, (short)0);
context.putShort(context.position() + HEADER_SIZE_LENGTH, (short)0);
}
writeElementAtOffset(context, HEADER_SIZE_LENGTH + (isDelta ? HEADER_ELT_LENGTH : 0), id, clock, value);
writeElementAtOffset(context, context.position() + HEADER_SIZE_LENGTH + (isDelta ? HEADER_ELT_LENGTH : 0), id, clock, value);
return context;
}
// write a tuple (node id, clock, count) at offset
// write a tuple (node id, clock, count) at an absolute (bytebuffer-wise) offset
private static void writeElementAtOffset(ByteBuffer context, int offset, NodeId id, long clock, long count)
{
ByteBufferUtil.arrayCopy(id.bytes(), id.bytes().position(), context, offset, NodeId.LENGTH);
context.putLong(offset + NodeId.LENGTH, clock);
context.putLong(offset + NodeId.LENGTH + CLOCK_LENGTH, count);
context = context.duplicate();
context.position(offset);
context.put(id.bytes().duplicate());
context.putLong(clock);
context.putLong(count);
}
private static int headerLength(ByteBuffer context)
@ -276,8 +282,9 @@ public class CounterContext implements IContext
*
* @param left counter context.
* @param right counter context.
* @param allocator An allocator for the merged value.
*/
public ByteBuffer merge(ByteBuffer left, ByteBuffer right)
public ByteBuffer merge(ByteBuffer left, ByteBuffer right, Allocator allocator)
{
ContextState leftState = new ContextState(left, headerLength(left));
ContextState rightState = new ContextState(right, headerLength(right));
@ -316,7 +323,7 @@ public class CounterContext implements IContext
mergedBodyLength += leftState.remainingBodyLength() + rightState.remainingBodyLength();
// Do the actual merge
ByteBuffer merged = ByteBuffer.allocate(mergedHeaderLength + mergedBodyLength);
ByteBuffer merged = allocator.allocate(mergedHeaderLength + mergedBodyLength);
merged.putShort(merged.position(), (short) ((mergedHeaderLength - HEADER_SIZE_LENGTH) / HEADER_ELT_LENGTH));
ContextState mergedState = new ContextState(merged, mergedHeaderLength);
leftState.reset();
@ -550,7 +557,7 @@ public class CounterContext implements IContext
// Found someone to merge it to
int nbDelta = foundState.isDelta() ? 1 : 0;
nbDelta += state.isDelta() ? 1 : 0;
ContextState merger = ContextState.allocate(2, nbDelta);
ContextState merger = ContextState.allocate(2, nbDelta, HeapAllocator.instance);
long fclock = foundState.getClock();
long fcount = foundState.getCount();
@ -644,7 +651,7 @@ public class CounterContext implements IContext
}
state.moveToNext();
}
return toAddBack == 0 ? cleanedContext : merge(cleanedContext, create(toAddBack));
return toAddBack == 0 ? cleanedContext : merge(cleanedContext, create(toAddBack, HeapAllocator.instance), HeapAllocator.instance);
}
/**
@ -792,11 +799,16 @@ public class CounterContext implements IContext
* ContextState corresponding.
*/
public static ContextState allocate(int elementCount, int deltaCount)
{
return allocate(elementCount, deltaCount, HeapAllocator.instance);
}
public static ContextState allocate(int elementCount, int deltaCount, Allocator allocator)
{
assert deltaCount <= elementCount;
int hlength = HEADER_SIZE_LENGTH + deltaCount * HEADER_ELT_LENGTH;
ByteBuffer context = ByteBuffer.allocate(hlength + elementCount * STEP_LENGTH);
context.putShort(0, (short)deltaCount);
ByteBuffer context = allocator.allocate(hlength + elementCount * STEP_LENGTH);
context.putShort(context.position(), (short)deltaCount);
return new ContextState(context, hlength);
}
}

View File

@ -19,6 +19,8 @@ package org.apache.cassandra.db.context;
import java.nio.ByteBuffer;
import org.apache.cassandra.utils.Allocator;
/**
* An opaque commutative context.
*
@ -53,12 +55,16 @@ public interface IContext
/**
* Return a context w/ an aggregated count for each node id.
*
* @param allocator
* An allocator for the merged value.
* @param left
* context.
* @param right
* context.
* @param allocator
* an allocator to allocate the new context from.
*/
public ByteBuffer merge(ByteBuffer left, ByteBuffer right);
public ByteBuffer merge(ByteBuffer left, ByteBuffer right, Allocator allocator);
/**
* Human-readable String from context.

View File

@ -44,6 +44,8 @@ import org.apache.cassandra.thrift.IndexClause;
import org.apache.cassandra.thrift.IndexExpression;
import org.apache.cassandra.thrift.IndexOperator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -248,7 +250,7 @@ public class KeysSearcher extends SecondaryIndexSearcher
assert !extraFilter.columns.isEmpty();
ColumnFamily cf = baseCfs.getColumnFamily(new QueryFilter(dk, path, extraFilter));
if (cf != null)
data.addAll(cf);
data.addAll(cf, HeapAllocator.instance);
}
}

View File

@ -25,6 +25,7 @@ import java.nio.ByteBuffer;
import org.apache.cassandra.db.*;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
public class CounterColumnType extends AbstractCommutativeType
{

View File

@ -23,10 +23,9 @@ package org.apache.cassandra.io;
import java.io.DataInput;
import java.io.IOException;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.IColumn;
public interface IColumnSerializer extends ICompactSerializer2<IColumn>
{
public IColumn deserialize(DataInput in, ColumnFamilyStore interner, boolean fromRemote, int expireBefore) throws IOException;
public IColumn deserialize(DataInput in, boolean fromRemote, int expireBefore) throws IOException;
}

View File

@ -175,7 +175,7 @@ public class SSTableIdentityIterator implements Comparable<SSTableIdentityIterat
{
try
{
IColumn column = columnFamily.getColumnSerializer().deserialize(inputWithTracker, null, fromRemote, expireBefore);
IColumn column = columnFamily.getColumnSerializer().deserialize(inputWithTracker, fromRemote, expireBefore);
if (validateColumns)
column.validateFields(columnFamily.metadata());
return column;
@ -230,7 +230,7 @@ public class SSTableIdentityIterator implements Comparable<SSTableIdentityIterat
assert inputWithTracker.getBytesRead() == headerSize();
ColumnFamily cf = columnFamily.cloneMeShallow();
// since we already read column count, just pass that value and continue deserialization
ColumnFamily.serializer().deserializeColumns(inputWithTracker, cf, columnCount, false, fromRemote);
ColumnFamily.serializer().deserializeColumns(inputWithTracker, cf, columnCount, fromRemote);
if (validateColumns)
{
try

View File

@ -28,6 +28,7 @@ import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.db.*;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.HeapAllocator;
/**
* A SSTable writer that doesn't assume rows are in sorted order.
@ -74,7 +75,7 @@ public class SSTableSimpleUnsortedWriter extends AbstractSSTableSimpleWriter
// Note that if the row was existing already, our size estimation will be slightly off
// since we'll be counting the key multiple times.
if (previous != null)
columnFamily.addAll(previous);
columnFamily.addAll(previous, HeapAllocator.instance);
if (currentSize > bufferSize)
sync();

View File

@ -43,15 +43,13 @@ public class ColumnSortedMap implements SortedMap<ByteBuffer, IColumn>
private final DataInput dis;
private final Comparator<ByteBuffer> comparator;
private final int length;
private final ColumnFamilyStore interner;
private final boolean fromRemote;
private final int expireBefore;
public ColumnSortedMap(Comparator<ByteBuffer> comparator, ColumnSerializer serializer, DataInput dis, ColumnFamilyStore interner, int length, boolean fromRemote, int expireBefore)
public ColumnSortedMap(Comparator<ByteBuffer> comparator, ColumnSerializer serializer, DataInput dis, int length, boolean fromRemote, int expireBefore)
{
this.comparator = comparator;
this.serializer = serializer;
this.interner = interner;
this.dis = dis;
this.length = length;
this.fromRemote = fromRemote;
@ -145,7 +143,7 @@ public class ColumnSortedMap implements SortedMap<ByteBuffer, IColumn>
public Set<Map.Entry<ByteBuffer, IColumn>> entrySet()
{
return new ColumnSet(serializer, dis, interner, length, fromRemote, expireBefore);
return new ColumnSet(serializer, dis, length, fromRemote, expireBefore);
}
}
@ -154,15 +152,13 @@ class ColumnSet implements Set<Map.Entry<ByteBuffer, IColumn>>
private final ColumnSerializer serializer;
private final DataInput dis;
private final int length;
private final ColumnFamilyStore interner;
private boolean fromRemote;
private final int expireBefore;
public ColumnSet(ColumnSerializer serializer, DataInput dis, ColumnFamilyStore interner, int length, boolean fromRemote, int expireBefore)
public ColumnSet(ColumnSerializer serializer, DataInput dis, int length, boolean fromRemote, int expireBefore)
{
this.serializer = serializer;
this.dis = dis;
this.interner = interner;
this.length = length;
this.fromRemote = fromRemote;
this.expireBefore = expireBefore;
@ -185,7 +181,7 @@ class ColumnSet implements Set<Map.Entry<ByteBuffer, IColumn>>
public Iterator<Entry<ByteBuffer, IColumn>> iterator()
{
return new ColumnIterator(serializer, dis, interner, length, fromRemote, expireBefore);
return new ColumnIterator(serializer, dis, length, fromRemote, expireBefore);
}
public Object[] toArray()
@ -240,14 +236,12 @@ class ColumnIterator implements Iterator<Map.Entry<ByteBuffer, IColumn>>
private final int length;
private final boolean fromRemote;
private int count = 0;
private ColumnFamilyStore interner;
private final int expireBefore;
public ColumnIterator(ColumnSerializer serializer, DataInput dis, ColumnFamilyStore interner, int length, boolean fromRemote, int expireBefore)
public ColumnIterator(ColumnSerializer serializer, DataInput dis, int length, boolean fromRemote, int expireBefore)
{
this.dis = dis;
this.serializer = serializer;
this.interner = interner;
this.length = length;
this.fromRemote = fromRemote;
this.expireBefore = expireBefore;
@ -258,7 +252,7 @@ class ColumnIterator implements Iterator<Map.Entry<ByteBuffer, IColumn>>
try
{
count++;
return serializer.deserialize(dis, interner, fromRemote, expireBefore);
return serializer.deserialize(dis, fromRemote, expireBefore);
}
catch (IOException e)
{

View File

@ -359,7 +359,7 @@ public class StorageProxy implements StorageProxyMBean
{
public void runMayThrow() throws IOException
{
rm.localCopy().apply();
rm.apply();
responseHandler.response(null);
}
};

View File

@ -148,7 +148,7 @@ public class IncomingStreamReader
// restore ColumnFamily
cf = ColumnFamily.create(cfs.metadata);
ColumnFamily.serializer().deserializeFromSSTableNoColumns(cf, in);
ColumnFamily.serializer().deserializeColumns(in, cf, true, true);
ColumnFamily.serializer().deserializeColumns(in, cf, true);
// write key and cf
writer.append(key, cf);

View File

@ -0,0 +1,41 @@
/**
* Copyright 2011 The Apache Software Foundation
*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.cassandra.utils;
import java.nio.ByteBuffer;
public abstract class Allocator
{
/**
* Allocate a slice of the given length.
*/
public ByteBuffer clone(ByteBuffer buffer)
{
assert buffer != null;
ByteBuffer cloned = allocate(buffer.remaining());
cloned.mark();
cloned.put(buffer.duplicate());
cloned.reset();
return cloned;
}
public abstract ByteBuffer allocate(int size);
}

View File

@ -169,11 +169,11 @@ public class ByteBufferUtil
if (buffer.hasArray())
{
int start = buffer.position();
if (buffer.arrayOffset() == 0 && start == 0 && length == buffer.array().length)
int boff = buffer.arrayOffset() + buffer.position();
if (boff == 0 && length == buffer.array().length)
return buffer.array();
else
return Arrays.copyOfRange(buffer.array(), start + buffer.arrayOffset(), start + length + buffer.arrayOffset());
return Arrays.copyOfRange(buffer.array(), boff, boff + length);
}
// else, DirectByteBuffer.get() is the fastest route
byte[] bytes = new byte[length];
@ -297,9 +297,8 @@ public class ByteBufferUtil
throw new IndexOutOfBoundsException();
for (int i = 0; i < length; i++)
{
// TODO: ByteBuffer.put is polymorphic, and might be slow here
dst.put(dstPos++, src.get(srcPos++));
}
}
}

View File

@ -0,0 +1,34 @@
/**
* Copyright 2011 The Apache Software Foundation
*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.cassandra.utils;
import java.nio.ByteBuffer;
public final class HeapAllocator extends Allocator
{
public static final HeapAllocator instance = new HeapAllocator();
private HeapAllocator() {}
public ByteBuffer allocate(int size)
{
return ByteBuffer.allocate(size);
}
}

View File

@ -0,0 +1,232 @@
/**
* Copyright 2011 The Apache Software Foundation
*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.cassandra.utils;
import java.nio.ByteBuffer;
import java.util.Collection;
import java.util.Collections;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import com.google.common.base.Preconditions;
import com.google.common.collect.Iterables;
/**
* The SlabAllocator is a bump-the-pointer allocator that allocates
* large (2MB by default) regions and then doles them out to threads that request
* slices into the array.
* <p/>
* The purpose of this class is to combat heap fragmentation in long lived
* objects: by ensuring that all allocations with similar lifetimes
* only to large regions of contiguous memory, we ensure that large blocks
* get freed up at the same time.
* <p/>
* Otherwise, variable length byte arrays allocated end up
* interleaved throughout the heap, and the old generation gets progressively
* more fragmented until a stop-the-world compacting collection occurs.
*/
public class SlabAllocator extends Allocator
{
private final static int REGION_SIZE = 2 * 1024 * 1024;
private final static int MAX_CLONED_SIZE = 256 * 1024; // bigger than this don't go in the region
private final AtomicReference<Region> currentRegion = new AtomicReference<Region>();
private final Collection<Region> filledRegions = new LinkedBlockingQueue<Region>();
/** @return Total number of bytes allocated by this allocator. */
public long size()
{
Iterable<Region> regions = filledRegions;
if (currentRegion.get() != null)
regions = Iterables.concat(regions, Collections.<Region>singleton(currentRegion.get()));
long total = 0;
for (Region region : regions)
total += region.size;
return total;
}
public ByteBuffer allocate(int size)
{
assert size >= 0;
if (size == 0)
return ByteBufferUtil.EMPTY_BYTE_BUFFER;
// satisfy large allocations directly from JVM since they don't cause fragmentation
// as badly, and fill up our regions quickly
if (size > MAX_CLONED_SIZE)
return ByteBuffer.allocate(size);
while (true)
{
Region region = getRegion();
// Try to allocate from this region
ByteBuffer cloned = region.allocate(size);
if (cloned != null)
return cloned;
// not enough space!
tryRetireRegion(region);
}
}
/**
* Try to retire the current region if it is still <code>region</code>.
* Postcondition is that curRegion.get() != region
*/
private void tryRetireRegion(Region region)
{
if (currentRegion.compareAndSet(region, null))
{
filledRegions.add(region);
}
}
/**
* Get the current region, or, if there is no current region, allocate a new one
*/
private Region getRegion()
{
while (true)
{
// Try to get the region
Region region = currentRegion.get();
if (region != null)
return region;
// No current region, so we want to allocate one. We race
// against other allocators to CAS in an uninitialized region
// (which is cheap to allocate)
region = new Region(REGION_SIZE);
if (currentRegion.compareAndSet(null, region))
{
// we won race - now we need to actually do the expensive allocation step
region.init();
return region;
}
// someone else won race - that's fine, we'll try to grab theirs
// in the next iteration of the loop.
}
}
/**
* A region of memory out of which allocations are sliced.
*
* This serves two purposes:
* - to provide a step between initialization and allocation, so that racing to CAS a
* new region in is harmless
* - encapsulates the allocation offset
*/
private static class Region
{
/**
* Actual underlying data
*/
private ByteBuffer data;
private static final int UNINITIALIZED = -1;
/**
* Offset for the next allocation, or the sentinel value -1
* which implies that the region is still uninitialized.
*/
private AtomicInteger nextFreeOffset = new AtomicInteger(UNINITIALIZED);
/**
* Total number of allocations satisfied from this buffer
*/
private AtomicInteger allocCount = new AtomicInteger();
/**
* Size of region in bytes
*/
private final int size;
/**
* Create an uninitialized region. Note that memory is not allocated yet, so
* this is cheap.
*
* @param size in bytes
*/
private Region(int size)
{
this.size = size;
}
/**
* Actually claim the memory for this region. This should only be called from
* the thread that constructed the region. It is thread-safe against other
* threads calling alloc(), who will block until the allocation is complete.
*/
public void init()
{
assert nextFreeOffset.get() == UNINITIALIZED;
data = ByteBuffer.allocate(size);
assert data.remaining() == data.capacity();
// Mark that it's ready for use
boolean initted = nextFreeOffset.compareAndSet(UNINITIALIZED, 0);
// We should always succeed the above CAS since only one thread calls init()!
Preconditions.checkState(initted, "Multiple threads tried to init same region");
}
/**
* Try to allocate <code>size</code> bytes from the region.
*
* @return the successful allocation, or null to indicate not-enough-space
*/
public ByteBuffer allocate(int size)
{
while (true)
{
int oldOffset = nextFreeOffset.get();
if (oldOffset == UNINITIALIZED)
{
// The region doesn't have its data allocated yet.
// Since we found this in currentRegion, we know that whoever
// CAS-ed it there is allocating it right now. So spin-loop
// shouldn't spin long!
Thread.yield();
continue;
}
if (oldOffset + size > data.capacity()) // capacity == remaining
return null;
// Try to atomically claim this region
if (nextFreeOffset.compareAndSet(oldOffset, oldOffset + size))
{
// we got the alloc
allocCount.incrementAndGet();
return (ByteBuffer) data.duplicate().position(oldOffset).limit(oldOffset + size);
}
// we raced and lost alloc, try again
}
}
@Override
public String toString()
{
return "Region@" + System.identityHashCode(this) +
" allocs=" + allocCount.get() + "waste=" +
(data.capacity() - nextFreeOffset.get());
}
}
}

View File

@ -1,21 +1,15 @@
package org.apache.cassandra.db;
import java.io.ByteArrayInputStream;
import java.io.DataInputStream;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.*;
import org.apache.cassandra.SchemaLoader;
import org.junit.Test;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.db.filter.QueryPath;
import static org.apache.cassandra.Util.column;
import static org.junit.Assert.*;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.db.marshal.BytesType;
import org.apache.cassandra.utils.HeapAllocator;
public class ArrayBackedSortedColumnsTest
{
@ -32,7 +26,7 @@ public class ArrayBackedSortedColumnsTest
int[] values = new int[]{ 1, 2, 2, 3 };
for (int i = 0; i < values.length; ++i)
map.addColumn(new Column(ByteBufferUtil.bytes(values[reversed ? values.length - 1 - i : i])));
map.addColumn(new Column(ByteBufferUtil.bytes(values[reversed ? values.length - 1 - i : i])), HeapAllocator.instance);
Iterator<IColumn> iter = map.iterator();
assertEquals("1st column", 1, iter.next().name().getInt(0));
@ -56,12 +50,12 @@ public class ArrayBackedSortedColumnsTest
int[] values2 = new int[]{ 2, 4, 5, 6 };
for (int i = 0; i < values1.length; ++i)
map.addColumn(new Column(ByteBufferUtil.bytes(values1[reversed ? values1.length - 1 - i : i])));
map.addColumn(new Column(ByteBufferUtil.bytes(values1[reversed ? values1.length - 1 - i : i])), HeapAllocator.instance);
for (int i = 0; i < values2.length; ++i)
map2.addColumn(new Column(ByteBufferUtil.bytes(values2[reversed ? values2.length - 1 - i : i])));
map2.addColumn(new Column(ByteBufferUtil.bytes(values2[reversed ? values2.length - 1 - i : i])), HeapAllocator.instance);
map2.addAll(map);
map2.addAll(map, HeapAllocator.instance);
Iterator<IColumn> iter = map2.iterator();
assertEquals("1st column", 1, iter.next().name().getInt(0));
@ -91,7 +85,7 @@ public class ArrayBackedSortedColumnsTest
Collections.reverse(reverseSorted);
for (int i = 0; i < values.length; ++i)
map.addColumn(new Column(ByteBufferUtil.bytes(values[reversed ? values.length - 1 - i : i])));
map.addColumn(new Column(ByteBufferUtil.bytes(values[reversed ? values.length - 1 - i : i])), HeapAllocator.instance);
assertSame(sorted, map.getSortedColumns());
assertSame(reverseSorted, map.getReverseSortedColumns());
@ -113,7 +107,7 @@ public class ArrayBackedSortedColumnsTest
names.add(ByteBufferUtil.bytes(v));
for (int i = 0; i < values.length; ++i)
map.addColumn(new Column(ByteBufferUtil.bytes(values[reversed ? values.length - 1 - i : i])));
map.addColumn(new Column(ByteBufferUtil.bytes(values[reversed ? values.length - 1 - i : i])), HeapAllocator.instance);
assertSame(names, map.getColumnNames());
}

View File

@ -31,6 +31,7 @@ import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.db.filter.QueryPath;
import static org.apache.cassandra.Util.column;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.HeapAllocator;
public class ColumnFamilyTest extends SchemaLoader
@ -126,8 +127,8 @@ public class ColumnFamilyTest extends SchemaLoader
cf_old.addColumn(QueryPath.column(ByteBufferUtil.bytes("col2")), val2, 1);
cf_old.addColumn(QueryPath.column(ByteBufferUtil.bytes("col3")), val2, 2);
cf_result.addAll(cf_new);
cf_result.addAll(cf_old);
cf_result.addAll(cf_new, HeapAllocator.instance);
cf_result.addAll(cf_old, HeapAllocator.instance);
assert 3 == cf_result.getColumnCount() : "Count is " + cf_new.getColumnCount();
//addcolumns will only add if timestamp >= old timestamp

View File

@ -40,6 +40,8 @@ import org.apache.cassandra.db.context.CounterContext;
import static org.apache.cassandra.db.context.CounterContext.ContextState;
import org.apache.cassandra.db.marshal.*;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.utils.Allocator;
import org.apache.cassandra.utils.HeapAllocator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.NodeId;
@ -205,6 +207,7 @@ public class CounterColumnTest extends SchemaLoader
@Test
public void testDiff() throws UnknownHostException
{
Allocator allocator = HeapAllocator.instance;
ContextState left;
ContextState right;
@ -226,7 +229,7 @@ public class CounterColumnTest extends SchemaLoader
assert null == rightCol.diff(leftCol);
// equality: equal nodes, all counts same
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 3L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -237,13 +240,13 @@ public class CounterColumnTest extends SchemaLoader
assert null == leftCol.diff(rightCol);
// greater than: left has superset of nodes (counts equal)
left = ContextState.allocate(4, 0);
left = ContextState.allocate(4, 0, allocator);
left.writeElement(NodeId.fromInt(3), 3L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
left.writeElement(NodeId.fromInt(12), 0L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 3L, 0L);
right.writeElement(NodeId.fromInt(6), 2L, 0L);
right.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -256,12 +259,12 @@ public class CounterColumnTest extends SchemaLoader
assert leftCol == rightCol.diff(leftCol);
// disjoint: right and left have disjoint node sets
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 1L, 0L);
left.writeElement(NodeId.fromInt(4), 1L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 1L, 0L);
right.writeElement(NodeId.fromInt(6), 1L, 0L);
right.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -275,7 +278,8 @@ public class CounterColumnTest extends SchemaLoader
@Test
public void testSerializeDeserialize() throws IOException
{
CounterContext.ContextState state = CounterContext.ContextState.allocate(4, 2);
Allocator allocator = HeapAllocator.instance;
CounterContext.ContextState state = CounterContext.ContextState.allocate(4, 2, allocator);
state.writeElement(NodeId.fromInt(1), 4L, 4L);
state.writeElement(NodeId.fromInt(2), 4L, 4L, true);
state.writeElement(NodeId.fromInt(3), 4L, 4L);
@ -291,7 +295,7 @@ public class CounterColumnTest extends SchemaLoader
assert original.equals(deserialized);
bufIn = new ByteArrayInputStream(serialized, 0, serialized.length);
CounterColumn deserializedOnRemote = (CounterColumn)Column.serializer().deserialize(new DataInputStream(bufIn), null, true);
CounterColumn deserializedOnRemote = (CounterColumn)Column.serializer().deserialize(new DataInputStream(bufIn), true);
assert deserializedOnRemote.name().equals(original.name());
assert deserializedOnRemote.total() == original.total();
assert deserializedOnRemote.value().equals(cc.clearAllDelta(original.value()));
@ -302,10 +306,11 @@ public class CounterColumnTest extends SchemaLoader
@Test
public void testUpdateDigest() throws Exception
{
Allocator allocator = HeapAllocator.instance;
MessageDigest digest1 = MessageDigest.getInstance("md5");
MessageDigest digest2 = MessageDigest.getInstance("md5");
CounterContext.ContextState state = CounterContext.ContextState.allocate(4, 2);
CounterContext.ContextState state = CounterContext.ContextState.allocate(4, 2, allocator);
state.writeElement(NodeId.fromInt(1), 4L, 4L);
state.writeElement(NodeId.fromInt(2), 4L, 4L, true);
state.writeElement(NodeId.fromInt(3), 4L, 4L);

View File

@ -164,7 +164,7 @@ public class CounterMutationTest extends CleanupHelper
}
// Check that if we merge old and clean on another node, we keep the right count
ByteBuffer onRemote = ctx.merge(ctx.clearAllDelta(state.context), ctx.clearAllDelta(cleaned));
ByteBuffer onRemote = ctx.merge(ctx.clearAllDelta(state.context), ctx.clearAllDelta(cleaned), HeapAllocator.instance);
assert ctx.total(onRemote) == 11;
}
}

View File

@ -95,7 +95,7 @@ public class ReadMessageTest extends SchemaLoader
ReadCommand command = new SliceByNamesReadCommand("Keyspace1", dk.key, new QueryPath("Standard1"), Arrays.asList(ByteBufferUtil.bytes("Column1")));
Row row = command.getRow(table);
IColumn col = row.cf.getColumn(ByteBufferUtil.bytes("Column1"));
assert Arrays.equals(col.value().array(), "abcd".getBytes());
assertEquals(col.value(), ByteBuffer.wrap("abcd".getBytes()));
}
@Test

View File

@ -110,8 +110,8 @@ public class TimeSortTest extends CleanupHelper
columnNames.add(getBytes(10));
columnNames.add(getBytes(0));
cf = cfStore.getColumnFamily(QueryFilter.getNamesFilter(Util.dk("900"), new QueryPath("StandardLong1"), columnNames));
assert "c".equals(new String(cf.getColumn(getBytes(0)).value().array()));
assert "c".equals(new String(cf.getColumn(getBytes(10)).value().array()));
assert "c".equals(ByteBufferUtil.string(cf.getColumn(getBytes(0)).value()));
assert "c".equals(ByteBufferUtil.string(cf.getColumn(getBytes(10)).value()));
}
private void validateTimeSort(Table table) throws IOException

View File

@ -25,17 +25,15 @@ import static org.junit.Assert.*;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.nio.ByteBuffer;
import java.util.*;
import org.apache.commons.lang.ArrayUtils;
import java.util.ArrayList;
import java.util.List;
import org.junit.Test;
import org.apache.cassandra.Util;
import org.apache.cassandra.db.context.IContext.ContextRelationship;
import static org.apache.cassandra.db.context.CounterContext.ContextState;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.NodeId;
import org.apache.cassandra.utils.*;
public class CounterContextTest
{
@ -55,17 +53,37 @@ public class CounterContextTest
stepLength = idLength + clockLength + countLength;
}
/** Allocates 1 byte from a new SlabAllocator and returns it. */
private Allocator bumpedSlab()
{
SlabAllocator allocator = new SlabAllocator();
allocator.allocate(1);
return allocator;
}
@Test
public void testCreate()
{
ByteBuffer context = cc.create(4);
assert context.remaining() == stepLength + 4;
runCreate(HeapAllocator.instance);
runCreate(bumpedSlab());
}
private void runCreate(Allocator allocator)
{
ByteBuffer bytes = cc.create(4, allocator);
assertEquals(stepLength + 4, bytes.remaining());
}
@Test
public void testDiff()
{
ContextState left = ContextState.allocate(3, 0);
runDiff(HeapAllocator.instance);
runDiff(bumpedSlab());
}
private void runDiff(Allocator allocator)
{
ContextState left = ContextState.allocate(3, 0, allocator);
ContextState right;
// equality: equal nodes, all counts same
@ -78,13 +96,13 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// greater than: left has superset of nodes (counts equal)
left = ContextState.allocate(4, 0);
left = ContextState.allocate(4, 0, allocator);
left.writeElement(NodeId.fromInt(3), 3L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
left.writeElement(NodeId.fromInt(12), 0L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 3L, 0L);
right.writeElement(NodeId.fromInt(6), 2L, 0L);
right.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -93,12 +111,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// less than: left has subset of nodes (counts equal)
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 3L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
right = ContextState.allocate(4, 0);
right = ContextState.allocate(4, 0, allocator);
right.writeElement(NodeId.fromInt(3), 3L, 0L);
right.writeElement(NodeId.fromInt(6), 2L, 0L);
right.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -108,12 +126,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// greater than: equal nodes, but left has higher counts
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 3L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 3L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 3L, 0L);
right.writeElement(NodeId.fromInt(6), 2L, 0L);
right.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -122,12 +140,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// less than: equal nodes, but right has higher counts
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 3L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 3L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 3L, 0L);
right.writeElement(NodeId.fromInt(6), 9L, 0L);
right.writeElement(NodeId.fromInt(9), 3L, 0L);
@ -136,12 +154,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// disjoint: right and left have disjoint node sets
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 1L, 0L);
left.writeElement(NodeId.fromInt(4), 1L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 1L, 0L);
right.writeElement(NodeId.fromInt(6), 1L, 0L);
right.writeElement(NodeId.fromInt(9), 1L, 0L);
@ -149,12 +167,12 @@ public class CounterContextTest
assert ContextRelationship.DISJOINT ==
cc.diff(left.context, right.context);
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 1L, 0L);
left.writeElement(NodeId.fromInt(4), 1L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(2), 1L, 0L);
right.writeElement(NodeId.fromInt(6), 1L, 0L);
right.writeElement(NodeId.fromInt(12), 1L, 0L);
@ -163,12 +181,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// disjoint: equal nodes, but right and left have higher counts in differing nodes
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 1L, 0L);
left.writeElement(NodeId.fromInt(6), 3L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 1L, 0L);
right.writeElement(NodeId.fromInt(6), 1L, 0L);
right.writeElement(NodeId.fromInt(9), 5L, 0L);
@ -176,12 +194,12 @@ public class CounterContextTest
assert ContextRelationship.DISJOINT ==
cc.diff(left.context, right.context);
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 2L, 0L);
left.writeElement(NodeId.fromInt(6), 3L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 1L, 0L);
right.writeElement(NodeId.fromInt(6), 9L, 0L);
right.writeElement(NodeId.fromInt(9), 5L, 0L);
@ -190,13 +208,13 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// disjoint: left has more nodes, but lower counts
left = ContextState.allocate(4, 0);
left = ContextState.allocate(4, 0, allocator);
left.writeElement(NodeId.fromInt(3), 2L, 0L);
left.writeElement(NodeId.fromInt(6), 3L, 0L);
left.writeElement(NodeId.fromInt(9), 1L, 0L);
left.writeElement(NodeId.fromInt(12), 1L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 4L, 0L);
right.writeElement(NodeId.fromInt(6), 9L, 0L);
right.writeElement(NodeId.fromInt(9), 5L, 0L);
@ -205,12 +223,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// disjoint: left has less nodes, but higher counts
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 5L, 0L);
left.writeElement(NodeId.fromInt(6), 3L, 0L);
left.writeElement(NodeId.fromInt(9), 2L, 0L);
right = ContextState.allocate(4, 0);
right = ContextState.allocate(4, 0, allocator);
right.writeElement(NodeId.fromInt(3), 4L, 0L);
right.writeElement(NodeId.fromInt(6), 3L, 0L);
right.writeElement(NodeId.fromInt(9), 2L, 0L);
@ -220,12 +238,12 @@ public class CounterContextTest
cc.diff(left.context, right.context);
// disjoint: mixed nodes and counts
left = ContextState.allocate(3, 0);
left = ContextState.allocate(3, 0, allocator);
left.writeElement(NodeId.fromInt(3), 5L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 2L, 0L);
right = ContextState.allocate(4, 0);
right = ContextState.allocate(4, 0, allocator);
right.writeElement(NodeId.fromInt(3), 4L, 0L);
right.writeElement(NodeId.fromInt(6), 3L, 0L);
right.writeElement(NodeId.fromInt(9), 2L, 0L);
@ -234,13 +252,13 @@ public class CounterContextTest
assert ContextRelationship.DISJOINT ==
cc.diff(left.context, right.context);
left = ContextState.allocate(4, 0);
left = ContextState.allocate(4, 0, allocator);
left.writeElement(NodeId.fromInt(3), 5L, 0L);
left.writeElement(NodeId.fromInt(6), 2L, 0L);
left.writeElement(NodeId.fromInt(7), 2L, 0L);
left.writeElement(NodeId.fromInt(9), 2L, 0L);
right = ContextState.allocate(3, 0);
right = ContextState.allocate(3, 0, allocator);
right.writeElement(NodeId.fromInt(3), 4L, 0L);
right.writeElement(NodeId.fromInt(6), 3L, 0L);
right.writeElement(NodeId.fromInt(9), 2L, 0L);
@ -251,61 +269,73 @@ public class CounterContextTest
@Test
public void testMerge()
{
runMerge(HeapAllocator.instance);
runMerge(bumpedSlab());
}
private void runMerge(Allocator allocator)
{
// note: local counts aggregated; remote counts are reconciled (i.e. take max)
ContextState left = ContextState.allocate(4, 1);
ContextState left = ContextState.allocate(4, 1, allocator);
left.writeElement(NodeId.fromInt(1), 1L, 1L);
left.writeElement(NodeId.fromInt(2), 2L, 2L);
left.writeElement(NodeId.fromInt(4), 6L, 3L);
left.writeElement(NodeId.getLocalId(), 7L, 3L, true);
ContextState right = ContextState.allocate(3, 1);
ContextState right = ContextState.allocate(3, 1, allocator);
right.writeElement(NodeId.fromInt(4), 4L, 4L);
right.writeElement(NodeId.fromInt(5), 5L, 5L);
right.writeElement(NodeId.getLocalId(), 2L, 9L, true);
ByteBuffer merged = cc.merge(left.context, right.context);
ByteBuffer merged = cc.merge(left.context, right.context, allocator);
int hd = 4;
assertEquals(hd + 5 * stepLength, merged.remaining());
// local node id's counts are aggregated
assert Util.equalsNodeId(NodeId.getLocalId(), merged, hd + 4*stepLength);
assertEquals( 9L, merged.getLong(hd + 4*stepLength + idLength));
assertEquals(12L, merged.getLong(hd + 4*stepLength + idLength + clockLength));
assertEquals( 9L, merged.getLong(merged.position() + hd + 4*stepLength + idLength));
assertEquals(12L, merged.getLong(merged.position() + hd + 4*stepLength + idLength + clockLength));
// remote node id counts are reconciled (i.e. take max)
assert Util.equalsNodeId(NodeId.fromInt(4), merged, hd + 2*stepLength);
assertEquals( 6L, merged.getLong(hd + 2*stepLength + idLength));
assertEquals( 3L, merged.getLong(hd + 2*stepLength + idLength + clockLength));
assertEquals( 6L, merged.getLong(merged.position() + hd + 2*stepLength + idLength));
assertEquals( 3L, merged.getLong(merged.position() + hd + 2*stepLength + idLength + clockLength));
assert Util.equalsNodeId(NodeId.fromInt(5), merged, hd + 3*stepLength);
assertEquals( 5L, merged.getLong(hd + 3*stepLength + idLength));
assertEquals( 5L, merged.getLong(hd + 3*stepLength + idLength + clockLength));
assertEquals( 5L, merged.getLong(merged.position() + hd + 3*stepLength + idLength));
assertEquals( 5L, merged.getLong(merged.position() + hd + 3*stepLength + idLength + clockLength));
assert Util.equalsNodeId(NodeId.fromInt(2), merged, hd + 1*stepLength);
assertEquals( 2L, merged.getLong(hd + 1*stepLength + idLength));
assertEquals( 2L, merged.getLong(hd + 1*stepLength + idLength + clockLength));
assertEquals( 2L, merged.getLong(merged.position() + hd + 1*stepLength + idLength));
assertEquals( 2L, merged.getLong(merged.position() + hd + 1*stepLength + idLength + clockLength));
assert Util.equalsNodeId(NodeId.fromInt(1), merged, hd + 0*stepLength);
assertEquals( 1L, merged.getLong(hd + 0*stepLength + idLength));
assertEquals( 1L, merged.getLong(hd + 0*stepLength + idLength + clockLength));
assertEquals( 1L, merged.getLong(merged.position() + hd + 0*stepLength + idLength));
assertEquals( 1L, merged.getLong(merged.position() + hd + 0*stepLength + idLength + clockLength));
}
@Test
public void testTotal()
{
ContextState left = ContextState.allocate(4, 1);
runTotal(HeapAllocator.instance);
runTotal(bumpedSlab());
}
private void runTotal(Allocator allocator)
{
ContextState left = ContextState.allocate(4, 1, allocator);
left.writeElement(NodeId.fromInt(1), 1L, 1L);
left.writeElement(NodeId.fromInt(2), 2L, 2L);
left.writeElement(NodeId.fromInt(4), 3L, 3L);
left.writeElement(NodeId.getLocalId(), 3L, 3L, true);
ContextState right = ContextState.allocate(3, 1);
ContextState right = ContextState.allocate(3, 1, allocator);
right.writeElement(NodeId.fromInt(4), 4L, 4L);
right.writeElement(NodeId.fromInt(5), 5L, 5L);
right.writeElement(NodeId.getLocalId(), 9L, 9L, true);
ByteBuffer merged = cc.merge(left.context, right.context);
ByteBuffer merged = cc.merge(left.context, right.context, allocator);
// 127.0.0.1: 12 (3+9)
// 0.0.0.1: 1
@ -318,6 +348,12 @@ public class CounterContextTest
@Test
public void testMergeOldShards()
{
runMergeOldShards(HeapAllocator.instance);
runMergeOldShards(bumpedSlab());
}
private void runMergeOldShards(Allocator allocator)
{
long now = System.currentTimeMillis();
NodeId id1 = NodeId.fromInt(1);
@ -327,7 +363,7 @@ public class CounterContextTest
records.add(new NodeId.NodeIdRecord(id3, 4L));
// Destination of merge is a delta
ContextState ctx = ContextState.allocate(5, 2);
ContextState ctx = ContextState.allocate(5, 2, allocator);
ctx.writeElement(id1, 1L, 1L);
ctx.writeElement(NodeId.fromInt(2), 2L, 2L);
ctx.writeElement(id3, 3L, 3L, true);
@ -344,41 +380,47 @@ public class CounterContextTest
assert m.getNodeId().equals(id3);
assert m.getClock() == 4L;
assert m.getCount() == 1L;
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger));
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger, allocator));
// Source of merge is a delta
ctx = ContextState.allocate(4, 1);
ctx = ContextState.allocate(4, 1, allocator);
ctx.writeElement(id1, 1L, 1L, true);
ctx.writeElement(NodeId.fromInt(2), 2L, 2L);
ctx.writeElement(id3, 3L, 3L);
ctx.writeElement(NodeId.fromInt(4), 6L, 3L);
merger = cc.computeOldShardMerger(ctx.context, records);
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger));
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger, allocator));
// source and destination of merge are deltas
ctx = ContextState.allocate(4, 2);
ctx = ContextState.allocate(4, 2, allocator);
ctx.writeElement(id1, 1L, 1L, true);
ctx.writeElement(NodeId.fromInt(2), 2L, 2L);
ctx.writeElement(id3, 3L, 3L, true);
ctx.writeElement(NodeId.fromInt(4), 6L, 3L);
merger = cc.computeOldShardMerger(ctx.context, records);
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger));
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger, allocator));
// none of source and destination of merge are deltas
ctx = ContextState.allocate(4, 0);
ctx = ContextState.allocate(4, 0, allocator);
ctx.writeElement(id1, 1L, 1L);
ctx.writeElement(NodeId.fromInt(2), 2L, 2L);
ctx.writeElement(id3, 3L, 3L);
ctx.writeElement(NodeId.fromInt(4), 6L, 3L);
merger = cc.computeOldShardMerger(ctx.context, records);
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger));
assert cc.total(ctx.context) == cc.total(cc.merge(ctx.context, merger, allocator));
}
@Test
public void testRemoveOldShards()
{
runRemoveOldShards(HeapAllocator.instance);
runRemoveOldShards(bumpedSlab());
}
private void runRemoveOldShards(Allocator allocator)
{
NodeId id1 = NodeId.fromInt(1);
NodeId id3 = NodeId.fromInt(3);
@ -388,7 +430,7 @@ public class CounterContextTest
records.add(new NodeId.NodeIdRecord(id3, 4L));
records.add(new NodeId.NodeIdRecord(id6, 10L));
ContextState ctx = ContextState.allocate(6, 2);
ContextState ctx = ContextState.allocate(6, 2, allocator);
ctx.writeElement(id1, 1L, 1L);
ctx.writeElement(NodeId.fromInt(2), 2L, 2L);
ctx.writeElement(id3, 3L, 3L, true);
@ -397,7 +439,7 @@ public class CounterContextTest
ctx.writeElement(id6, 5L, 6L);
ByteBuffer merger = cc.computeOldShardMerger(ctx.context, records);
ByteBuffer merged = cc.merge(ctx.context, merger);
ByteBuffer merged = cc.merge(ctx.context, merger, allocator);
assert cc.total(ctx.context) == cc.total(merged);
ByteBuffer cleaned = cc.removeOldShards(merged, (int)(System.currentTimeMillis() / 1000) + 1);
@ -405,12 +447,11 @@ public class CounterContextTest
assert cleaned.remaining() == ctx.context.remaining() - stepLength;
merger = cc.computeOldShardMerger(cleaned, records);
merged = cc.merge(cleaned, merger);
merged = cc.merge(cleaned, merger, allocator);
assert cc.total(ctx.context) == cc.total(merged);
cleaned = cc.removeOldShards(merged, (int)(System.currentTimeMillis() / 1000) + 1);
assert cc.total(ctx.context) == cc.total(cleaned);
assert cleaned.remaining() == ctx.context.remaining() - 2 * stepLength - 2;
}
}