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();
 


Reply via email to