Author: jbellis
Date: Wed Jun 16 15:26:35 2010
New Revision: 955263
URL: http://svn.apache.org/viewvc?rev=955263&view=rev
Log:
introduce AbstractCompactedRow, PrecompactedRow
patch by jbellis; reviewed by Stu Hood for CASSANDRA-16
Added:
cassandra/trunk/src/java/org/apache/cassandra/io/AbstractCompactedRow.java
cassandra/trunk/src/java/org/apache/cassandra/io/PrecompactedRow.java
Modified:
cassandra/trunk/src/java/org/apache/cassandra/db/CompactionManager.java
cassandra/trunk/src/java/org/apache/cassandra/io/CompactionIterator.java
cassandra/trunk/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java
cassandra/trunk/src/java/org/apache/cassandra/service/AntiEntropyService.java
cassandra/trunk/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java
Modified:
cassandra/trunk/src/java/org/apache/cassandra/db/CompactionManager.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/src/java/org/apache/cassandra/db/CompactionManager.java?rev=955263&r1=955262&r2=955263&view=diff
==============================================================================
--- cassandra/trunk/src/java/org/apache/cassandra/db/CompactionManager.java
(original)
+++ cassandra/trunk/src/java/org/apache/cassandra/db/CompactionManager.java Wed
Jun 16 15:26:35 2010
@@ -324,7 +324,7 @@ public class CompactionManager implement
SSTableWriter writer;
CompactionIterator ci = new CompactionIterator(sstables, gcBefore,
major); // retain a handle so we can call close()
- Iterator<CompactionIterator.CompactedRow> nni = new FilterIterator(ci,
PredicateUtils.notNullPredicate());
+ Iterator<AbstractCompactedRow> nni = new FilterIterator(ci,
PredicateUtils.notNullPredicate());
executor.beginCompaction(cfs, ci);
try
@@ -346,10 +346,10 @@ public class CompactionManager implement
validator.prepare();
while (nni.hasNext())
{
- CompactionIterator.CompactedRow row = nni.next();
+ AbstractCompactedRow row = nni.next();
long prevpos = writer.getFilePointer();
- writer.append(row.key, row.buffer);
+ writer.append(row);
validator.add(row);
totalkeysWritten++;
@@ -420,7 +420,7 @@ public class CompactionManager implement
SSTableWriter writer = null;
CompactionIterator ci = new AntiCompactionIterator(sstables, ranges,
getDefaultGCBefore(), cfs.isCompleteSSTables(sstables));
- Iterator<CompactionIterator.CompactedRow> nni = new FilterIterator(ci,
PredicateUtils.notNullPredicate());
+ Iterator<AbstractCompactedRow> nni = new FilterIterator(ci,
PredicateUtils.notNullPredicate());
executor.beginCompaction(cfs, ci);
try
@@ -432,14 +432,14 @@ public class CompactionManager implement
while (nni.hasNext())
{
- CompactionIterator.CompactedRow row = nni.next();
+ AbstractCompactedRow row = nni.next();
if (writer == null)
{
FileUtils.createDirectory(compactionFileLocation);
String newFilename = new
File(cfs.getTempSSTablePath(compactionFileLocation)).getAbsolutePath();
writer = new SSTableWriter(newFilename,
expectedBloomFilterSize, StorageService.getPartitioner());
}
- writer.append(row.key, row.buffer);
+ writer.append(row);
totalkeysWritten++;
}
}
@@ -486,14 +486,14 @@ public class CompactionManager implement
executor.beginCompaction(cfs, ci);
try
{
- Iterator<CompactionIterator.CompactedRow> nni = new
FilterIterator(ci, PredicateUtils.notNullPredicate());
+ Iterator<AbstractCompactedRow> nni = new FilterIterator(ci,
PredicateUtils.notNullPredicate());
// validate the CF as we iterate over it
AntiEntropyService.IValidator validator =
AntiEntropyService.instance.getValidator(cfs.getTable().name,
cfs.getColumnFamilyName(), initiator, true);
validator.prepare();
while (nni.hasNext())
{
- CompactionIterator.CompactedRow row = nni.next();
+ AbstractCompactedRow row = nni.next();
validator.add(row);
}
validator.complete();
Added:
cassandra/trunk/src/java/org/apache/cassandra/io/AbstractCompactedRow.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/src/java/org/apache/cassandra/io/AbstractCompactedRow.java?rev=955263&view=auto
==============================================================================
--- cassandra/trunk/src/java/org/apache/cassandra/io/AbstractCompactedRow.java
(added)
+++ cassandra/trunk/src/java/org/apache/cassandra/io/AbstractCompactedRow.java
Wed Jun 16 15:26:35 2010
@@ -0,0 +1,28 @@
+package org.apache.cassandra.io;
+
+import java.io.DataOutput;
+import java.io.IOException;
+import java.security.MessageDigest;
+
+import org.apache.cassandra.db.DecoratedKey;
+
+/**
+ * a CompactedRow is an object that takes a bunch of rows (keys +
columnfamilies)
+ * and can write a compacted version of those rows to an output stream. It
does
+ * NOT necessarily require creating a merged CF object in memory.
+ */
+public abstract class AbstractCompactedRow
+{
+ public final DecoratedKey key;
+
+ public AbstractCompactedRow(DecoratedKey key)
+ {
+ this.key = key;
+ }
+
+ public abstract void write(DataOutput out) throws IOException;
+
+ public abstract void update(MessageDigest digest);
+
+ public abstract boolean isEmpty();
+}
Modified:
cassandra/trunk/src/java/org/apache/cassandra/io/CompactionIterator.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/src/java/org/apache/cassandra/io/CompactionIterator.java?rev=955263&r1=955262&r2=955263&view=diff
==============================================================================
--- cassandra/trunk/src/java/org/apache/cassandra/io/CompactionIterator.java
(original)
+++ cassandra/trunk/src/java/org/apache/cassandra/io/CompactionIterator.java
Wed Jun 16 15:26:35 2010
@@ -42,7 +42,7 @@ import org.apache.cassandra.io.sstable.S
import org.apache.cassandra.io.sstable.SSTableScanner;
import org.apache.cassandra.io.util.DataOutputBuffer;
-public class CompactionIterator extends
ReducingIterator<SSTableIdentityIterator, CompactionIterator.CompactedRow>
implements Closeable
+public class CompactionIterator extends
ReducingIterator<SSTableIdentityIterator, AbstractCompactedRow> implements
Closeable
{
private static Logger logger =
LoggerFactory.getLogger(CompactionIterator.class);
@@ -97,55 +97,14 @@ public class CompactionIterator extends
rows.add(current);
}
- protected CompactedRow getReduced()
+ protected AbstractCompactedRow getReduced()
{
assert rows.size() > 0;
- DataOutputBuffer buffer = new DataOutputBuffer();
- DecoratedKey key = rows.get(0).getKey();
try
{
- if (rows.size() > 1 || major)
- {
- ColumnFamily cf = null;
- for (SSTableIdentityIterator row : rows)
- {
- ColumnFamily thisCF;
- try
- {
- thisCF = row.getColumnFamily();
- }
- catch (IOException e)
- {
- logger.error("Skipping row " + key + " in " +
row.getPath(), e);
- continue;
- }
- if (cf == null)
- {
- cf = thisCF;
- }
- else
- {
- cf.addAll(thisCF);
- }
- }
- ColumnFamily cfPurged = major ?
ColumnFamilyStore.removeDeleted(cf, gcBefore) : cf;
- if (cfPurged == null)
- return null;
- ColumnFamily.serializer().serializeWithIndexes(cfPurged,
buffer);
- }
- else
- {
- assert rows.size() == 1;
- try
- {
- rows.get(0).echoData(buffer);
- }
- catch (IOException e)
- {
- throw new IOError(e);
- }
- }
+ PrecompactedRow compactedRow = new PrecompactedRow(rows, major,
gcBefore);
+ return compactedRow.isEmpty() ? null : compactedRow;
}
finally
{
@@ -159,7 +118,6 @@ public class CompactionIterator extends
}
}
}
- return new CompactedRow(key, buffer);
}
public void close() throws IOException
@@ -185,15 +143,4 @@ public class CompactionIterator extends
return bytesRead;
}
- public static class CompactedRow
- {
- public final DecoratedKey key;
- public final DataOutputBuffer buffer;
-
- public CompactedRow(DecoratedKey key, DataOutputBuffer buffer)
- {
- this.key = key;
- this.buffer = buffer;
- }
- }
}
Added: cassandra/trunk/src/java/org/apache/cassandra/io/PrecompactedRow.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/src/java/org/apache/cassandra/io/PrecompactedRow.java?rev=955263&view=auto
==============================================================================
--- cassandra/trunk/src/java/org/apache/cassandra/io/PrecompactedRow.java
(added)
+++ cassandra/trunk/src/java/org/apache/cassandra/io/PrecompactedRow.java Wed
Jun 16 15:26:35 2010
@@ -0,0 +1,95 @@
+package org.apache.cassandra.io;
+
+import java.io.DataOutput;
+import java.io.IOError;
+import java.io.IOException;
+import java.security.MessageDigest;
+import java.util.List;
+
+import org.apache.cassandra.db.ColumnFamily;
+import org.apache.cassandra.db.ColumnFamilyStore;
+import org.apache.cassandra.db.DecoratedKey;
+import org.apache.cassandra.io.sstable.SSTableIdentityIterator;
+import org.apache.cassandra.io.util.DataOutputBuffer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * PrecompactedRow merges its rows in its constructor in memory.
+ */
+public class PrecompactedRow extends AbstractCompactedRow
+{
+ private static Logger logger =
LoggerFactory.getLogger(PrecompactedRow.class);
+
+ private final DataOutputBuffer buffer;
+
+ public PrecompactedRow(DecoratedKey key, DataOutputBuffer buffer)
+ {
+ super(key);
+ this.buffer = buffer;
+ }
+
+ public PrecompactedRow(List<SSTableIdentityIterator> rows, boolean major,
int gcBefore)
+ {
+ super(rows.get(0).getKey());
+ buffer = new DataOutputBuffer();
+
+ if (rows.size() > 1 || major)
+ {
+ ColumnFamily cf = null;
+ for (SSTableIdentityIterator row : rows)
+ {
+ ColumnFamily thisCF;
+ try
+ {
+ thisCF = row.getColumnFamily();
+ }
+ catch (IOException e)
+ {
+ logger.error("Skipping row " + key + " in " +
row.getPath(), e);
+ continue;
+ }
+ if (cf == null)
+ {
+ cf = thisCF;
+ }
+ else
+ {
+ cf.addAll(thisCF);
+ }
+ }
+ ColumnFamily cfPurged = major ?
ColumnFamilyStore.removeDeleted(cf, gcBefore) : cf;
+ if (cfPurged == null)
+ return;
+ ColumnFamily.serializer().serializeWithIndexes(cfPurged, buffer);
+ }
+ else
+ {
+ assert rows.size() == 1;
+ try
+ {
+ rows.get(0).echoData(buffer);
+ }
+ catch (IOException e)
+ {
+ throw new IOError(e);
+ }
+ }
+ }
+
+ public void write(DataOutput out) throws IOException
+ {
+ out.writeInt(buffer.getLength());
+ out.write(buffer.getData(), 0, buffer.getLength());
+ }
+
+ public void update(MessageDigest digest)
+ {
+ digest.update(buffer.getData(), 0, buffer.getLength());
+ }
+
+ public boolean isEmpty()
+ {
+ return buffer.getLength() == 0;
+ }
+}
Modified:
cassandra/trunk/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java?rev=955263&r1=955262&r2=955263&view=diff
==============================================================================
--- cassandra/trunk/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java
(original)
+++ cassandra/trunk/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java
Wed Jun 16 15:26:35 2010
@@ -42,6 +42,7 @@ import java.io.FileOutputStream;
import java.io.IOError;
import java.io.IOException;
+import org.apache.cassandra.io.AbstractCompactedRow;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -110,6 +111,14 @@ public class SSTableWriter extends SSTab
dbuilder.addPotentialBoundary(dataPosition);
}
+ public void append(AbstractCompactedRow row) throws IOException
+ {
+ long currentPosition = beforeAppend(row.key);
+
FBUtilities.writeShortByteArray(partitioner.convertToDiskFormat(row.key),
dataFile);
+ row.write(dataFile);
+ afterAppend(row.key, currentPosition);
+ }
+
// TODO make this take a DataOutputStream and wrap the byte[] version to
combine them
public void append(DecoratedKey decoratedKey, DataOutputBuffer buffer)
throws IOException
{
Modified:
cassandra/trunk/src/java/org/apache/cassandra/service/AntiEntropyService.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/src/java/org/apache/cassandra/service/AntiEntropyService.java?rev=955263&r1=955262&r2=955263&view=diff
==============================================================================
---
cassandra/trunk/src/java/org/apache/cassandra/service/AntiEntropyService.java
(original)
+++
cassandra/trunk/src/java/org/apache/cassandra/service/AntiEntropyService.java
Wed Jun 16 15:26:35 2010
@@ -20,6 +20,8 @@ package org.apache.cassandra.service;
import java.io.*;
import java.net.InetAddress;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
import java.util.*;
import java.util.concurrent.*;
@@ -31,11 +33,9 @@ import org.apache.cassandra.db.Decorated
import org.apache.cassandra.db.Table;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
-import org.apache.cassandra.io.CompactionIterator.CompactedRow;
+import org.apache.cassandra.io.AbstractCompactedRow;
import org.apache.cassandra.io.ICompactSerializer;
-import org.apache.cassandra.io.sstable.SSTable;
import org.apache.cassandra.io.sstable.SSTableReader;
-import org.apache.cassandra.io.sstable.IndexSummary;
import org.apache.cassandra.streaming.StreamOut;
import org.apache.cassandra.net.IVerbHandler;
import org.apache.cassandra.net.Message;
@@ -347,7 +347,7 @@ public class AntiEntropyService
public static interface IValidator
{
public void prepare();
- public void add(CompactedRow row);
+ public void add(AbstractCompactedRow row);
public void complete();
}
@@ -440,7 +440,7 @@ public class AntiEntropyService
*
* @param row The row.
*/
- public void add(CompactedRow row)
+ public void add(AbstractCompactedRow row)
{
if (mintoken != null)
{
@@ -471,12 +471,21 @@ public class AntiEntropyService
range.addHash(rowHash(row));
}
- private MerkleTree.RowHash rowHash(CompactedRow row)
+ private MerkleTree.RowHash rowHash(AbstractCompactedRow row)
{
validated++;
// MerkleTree uses XOR internally, so we want lots of output bits
here
- byte[] rowhash = FBUtilities.hash("SHA-256", row.key.key,
row.buffer.getData());
- return new MerkleTree.RowHash(row.key.token, rowhash);
+ MessageDigest digest = null;
+ try
+ {
+ digest = MessageDigest.getInstance("SHA-256");
+ }
+ catch (NoSuchAlgorithmException e)
+ {
+ throw new AssertionError(e);
+ }
+ row.update(digest);
+ return new MerkleTree.RowHash(row.key.token, digest.digest());
}
/**
@@ -541,7 +550,7 @@ public class AntiEntropyService
/**
* Does nothing.
*/
- public void add(CompactedRow row)
+ public void add(AbstractCompactedRow row)
{
// noop
}
Modified:
cassandra/trunk/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java
URL:
http://svn.apache.org/viewvc/cassandra/trunk/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java?rev=955263&r1=955262&r2=955263&view=diff
==============================================================================
---
cassandra/trunk/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java
(original)
+++
cassandra/trunk/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java
Wed Jun 16 15:26:35 2010
@@ -29,7 +29,7 @@ import org.apache.cassandra.db.*;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Token;
-import org.apache.cassandra.io.CompactionIterator.CompactedRow;
+import org.apache.cassandra.io.PrecompactedRow;
import org.apache.cassandra.io.util.DataOutputBuffer;
import org.apache.cassandra.locator.AbstractReplicationStrategy;
import org.apache.cassandra.locator.TokenMetadata;
@@ -38,7 +38,6 @@ import org.apache.cassandra.utils.FBUtil
import org.apache.cassandra.utils.MerkleTree;
import org.apache.cassandra.CleanupHelper;
-import org.apache.cassandra.config.DatabaseDescriptorTest;
import org.apache.cassandra.Util;
import org.junit.After;
@@ -156,11 +155,11 @@ public class AntiEntropyServiceTest exte
validator.prepare();
// add a row with the minimum token
- validator.add(new CompactedRow(new DecoratedKey(min,
"nonsense!".getBytes(FBUtilities.UTF8)),
+ validator.add(new PrecompactedRow(new DecoratedKey(min,
"nonsense!".getBytes(FBUtilities.UTF8)),
new DataOutputBuffer()));
// and a row after it
- validator.add(new CompactedRow(new DecoratedKey(mid,
"inconceivable!".getBytes(FBUtilities.UTF8)),
+ validator.add(new PrecompactedRow(new DecoratedKey(mid,
"inconceivable!".getBytes(FBUtilities.UTF8)),
new DataOutputBuffer()));
validator.complete();