diff --git a/src/org/apache/cassandra/db/Column.java b/src/org/apache/cassandra/db/Column.java index 68153656d8..74d70bb2d7 100644 --- a/src/org/apache/cassandra/db/Column.java +++ b/src/org/apache/cassandra/db/Column.java @@ -22,6 +22,7 @@ import java.io.DataInputStream; import java.io.DataOutputStream; import java.io.IOException; import java.util.Collection; +import java.nio.ByteBuffer; import org.apache.commons.lang.ArrayUtils; @@ -192,6 +193,11 @@ public final class Column implements IColumn return stringBuilder.toString().getBytes(); } + public int getLocalDeletionTime() + { + assert isMarkedForDelete; + return ByteBuffer.wrap(value).getInt(); + } } class ColumnSerializer implements ICompactSerializer2 diff --git a/src/org/apache/cassandra/db/ColumnFamily.java b/src/org/apache/cassandra/db/ColumnFamily.java index ec0a80c2c0..107d426582 100644 --- a/src/org/apache/cassandra/db/ColumnFamily.java +++ b/src/org/apache/cassandra/db/ColumnFamily.java @@ -96,6 +96,7 @@ public final class ColumnFamily private transient ICompactSerializer2 columnSerializer_; private long markedForDeleteAt = Long.MIN_VALUE; + private int localDeletionTime = Integer.MIN_VALUE; private AtomicInteger size_ = new AtomicInteger(0); private EfficientBidiMap columns_; @@ -156,11 +157,18 @@ public final class ColumnFamily createColumnFactoryAndColumnSerializer(columnType); } + ColumnFamily cloneMeShallow() + { + ColumnFamily cf = new ColumnFamily(name_, type_); + cf.markedForDeleteAt = markedForDeleteAt; + cf.localDeletionTime = localDeletionTime; + return cf; + } + ColumnFamily cloneMe() { - ColumnFamily cf = new ColumnFamily(name_, type_); - cf.markedForDeleteAt = markedForDeleteAt; - cf.columns_ = columns_.cloneMe(); + ColumnFamily cf = cloneMeShallow(); + cf.columns_ = columns_.cloneMe(); return cf; } @@ -292,8 +300,9 @@ public final class ColumnFamily columns_.remove(columnName); } - void delete(long timestamp) + void delete(int localtime, long timestamp) { + localDeletionTime = localtime; markedForDeleteAt = timestamp; } @@ -413,10 +422,16 @@ public final class ColumnFamily return xorHash; } - public long getMarkedForDeleteAt() { + public long getMarkedForDeleteAt() + { return markedForDeleteAt; } + public int getLocalDeletionTime() + { + return localDeletionTime; + } + public String type() { return type_; @@ -452,15 +467,11 @@ public final class ColumnFamily { Collection columns = columnFamily.getAllColumns(); - /* write the column family id */ dos.writeUTF(columnFamily.name()); - /* write if this cf is marked for delete */ - dos.writeLong(columnFamily.getMarkedForDeleteAt()); + dos.writeInt(columnFamily.localDeletionTime); + dos.writeLong(columnFamily.markedForDeleteAt); - /* write the size is the number of columns */ dos.writeInt(columns.size()); - - /* write the column data */ for ( IColumn column : columns ) { columnFamily.getColumnSerializer().serialize(column, dos); @@ -475,7 +486,7 @@ public final class ColumnFamily { String name = dis.readUTF(); ColumnFamily cf = new ColumnFamily(name, DatabaseDescriptor.getColumnFamilyType(name)); - cf.delete(dis.readLong()); + cf.delete(dis.readInt(), dis.readLong()); return cf; } diff --git a/src/org/apache/cassandra/db/ColumnFamilyStore.java b/src/org/apache/cassandra/db/ColumnFamilyStore.java index 00001a66a0..6cb974fbf4 100644 --- a/src/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/org/apache/cassandra/db/ColumnFamilyStore.java @@ -23,10 +23,6 @@ import java.io.IOException; import java.lang.management.ManagementFactory; import javax.management.MBeanServer; import javax.management.ObjectName; -import javax.management.InstanceAlreadyExistsException; -import javax.management.MBeanRegistrationException; -import javax.management.NotCompliantMBeanException; -import javax.management.MalformedObjectNameException; import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; @@ -594,14 +590,15 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean // start from nothing so that we don't include potential deleted columns from the first instance ColumnFamily cf0 = columnFamilies.get(0); - ColumnFamily cf = new ColumnFamily(cf0.name(), cf0.type()); + ColumnFamily cf = cf0.cloneMeShallow(); // merge for (ColumnFamily cf2 : columnFamilies) { assert cf.name().equals(cf2.name()); cf.addColumns(cf2); - cf.delete(Math.max(cf.getMarkedForDeleteAt(), cf2.getMarkedForDeleteAt())); + cf.delete(Math.max(cf.getLocalDeletionTime(), cf2.getLocalDeletionTime()), + Math.max(cf.getMarkedForDeleteAt(), cf2.getMarkedForDeleteAt())); } return cf; } @@ -619,29 +616,60 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return removeDeleted(cf); } - static ColumnFamily removeDeleted(ColumnFamily cf) { - if (cf == null) { + static final int GC_GRACE_IN_SECONDS = 10 * 24 * 3600; // 10 days + + /* + This is complicated because we need to preserve deleted columns, supercolumns, and columnfamilies + until they have been deleted for at least GC_GRACE_IN_SECONDS. But, we do not need to preserve + their contents; just the object itself as a "tombstone" that can be used to repair other + replicas that do not know about the deletion. + */ + static ColumnFamily removeDeleted(ColumnFamily cf) + { + return removeDeleted(cf, (int)(System.currentTimeMillis() / 1000) - GC_GRACE_IN_SECONDS); + } + + static ColumnFamily removeDeleted(ColumnFamily cf, int gcBefore) + { + if (cf == null) return null; - } - for (String cname : new ArrayList(cf.getColumns().keySet())) { + + for (String cname : new ArrayList(cf.getColumns().keySet())) + { IColumn c = cf.getColumns().get(cname); - if (c instanceof SuperColumn) { - long min_timestamp = Math.max(c.getMarkedForDeleteAt(), cf.getMarkedForDeleteAt()); + if (c instanceof SuperColumn) + { + long minTimestamp = Math.max(c.getMarkedForDeleteAt(), cf.getMarkedForDeleteAt()); // don't operate directly on the supercolumn, it could be the one in the memtable cf.remove(cname); - IColumn sc = new SuperColumn(cname); - for (IColumn subColumn : c.getSubColumns()) { - if (!subColumn.isMarkedForDelete() && subColumn.timestamp() >= min_timestamp) { - sc.addColumn(subColumn.name(), subColumn); + SuperColumn sc = new SuperColumn(cname); + sc.markForDeleteAt(c.getLocalDeletionTime(), c.getMarkedForDeleteAt()); + for (IColumn subColumn : c.getSubColumns()) + { + if (subColumn.timestamp() >= minTimestamp) + { + if (!subColumn.isMarkedForDelete() || subColumn.getLocalDeletionTime() > gcBefore) + { + sc.addColumn(subColumn.name(), subColumn); + } } } - if (sc.getSubColumns().size() > 0) { + if (sc.getSubColumns().size() > 0 || sc.getLocalDeletionTime() > gcBefore) + { cf.addColumn(sc); } - } else if (c.isMarkedForDelete() || c.timestamp() < cf.getMarkedForDeleteAt()) { + } + else if ((c.isMarkedForDelete() && c.getLocalDeletionTime() <= gcBefore) + || c.timestamp() < cf.getMarkedForDeleteAt()) + { cf.remove(cname); } } + + if (cf.getColumnCount() == 0 && cf.getLocalDeletionTime() <= gcBefore) + { + return null; + } return cf; } diff --git a/src/org/apache/cassandra/db/IColumn.java b/src/org/apache/cassandra/db/IColumn.java index 1f0284cf3f..fc5784cafa 100644 --- a/src/org/apache/cassandra/db/IColumn.java +++ b/src/org/apache/cassandra/db/IColumn.java @@ -43,4 +43,5 @@ public interface IColumn public IColumn diff(IColumn column); public int getObjectCount(); public byte[] digest(); + public int getLocalDeletionTime(); // for tombstone GC, so int is sufficient granularity } diff --git a/src/org/apache/cassandra/db/Memtable.java b/src/org/apache/cassandra/db/Memtable.java index 33bb98a3bf..bc854748e5 100644 --- a/src/org/apache/cassandra/db/Memtable.java +++ b/src/org/apache/cassandra/db/Memtable.java @@ -276,7 +276,8 @@ public class Memtable implements Comparable int newObjectCount = oldCf.getColumnCount(); resolveSize(oldSize, newSize); resolveCount(oldObjectCount, newObjectCount); - oldCf.delete(Math.max(oldCf.getMarkedForDeleteAt(), columnFamily.getMarkedForDeleteAt())); + oldCf.delete(Math.max(oldCf.getLocalDeletionTime(), columnFamily.getLocalDeletionTime()), + Math.max(oldCf.getMarkedForDeleteAt(), columnFamily.getMarkedForDeleteAt())); } else { diff --git a/src/org/apache/cassandra/db/RowMutation.java b/src/org/apache/cassandra/db/RowMutation.java index 19cee2e62c..89ec0dc521 100644 --- a/src/org/apache/cassandra/db/RowMutation.java +++ b/src/org/apache/cassandra/db/RowMutation.java @@ -29,6 +29,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.ExecutionException; +import java.nio.ByteBuffer; import org.apache.commons.lang.ArrayUtils; import org.apache.commons.lang.StringUtils; @@ -175,14 +176,16 @@ public class RowMutation implements Serializable { String[] values = RowMutation.getColumnAndColumnFamily(columnFamilyColumn); String cfName = values[0]; + if (modifications_.containsKey(cfName)) { throw new IllegalArgumentException("ColumnFamily " + cfName + " is already being modified"); } - if (values.length == 0 || values.length > 3) throw new IllegalArgumentException("Column Family " + columnFamilyColumn + " in invalid format. Must be in : format."); + int localDeleteTime = (int) (System.currentTimeMillis() / 1000); + ColumnFamily columnFamily = modifications_.get(cfName); if (columnFamily == null) columnFamily = new ColumnFamily(cfName, DatabaseDescriptor.getColumnType(cfName)); @@ -191,22 +194,26 @@ public class RowMutation implements Serializable if (columnFamily.isSuper()) { SuperColumn sc = new SuperColumn(values[1]); - sc.markForDeleteAt(timestamp); + sc.markForDeleteAt(localDeleteTime, timestamp); columnFamily.addColumn(sc); } else { - columnFamily.addColumn(values[1], ArrayUtils.EMPTY_BYTE_ARRAY, timestamp, true); + ByteBuffer bytes = ByteBuffer.allocate(4); + bytes.putInt(localDeleteTime); + columnFamily.addColumn(values[1], bytes.array(), timestamp, true); } } else if (values.length == 3) { - columnFamily.addColumn(values[1] + ":" + values[2], ArrayUtils.EMPTY_BYTE_ARRAY, timestamp, true); + ByteBuffer bytes = ByteBuffer.allocate(4); + bytes.putInt(localDeleteTime); + columnFamily.addColumn(values[1] + ":" + values[2], bytes.array(), timestamp, true); } else { assert values.length == 1; - columnFamily.delete(timestamp); + columnFamily.delete(localDeleteTime, timestamp); } modifications_.put(cfName, columnFamily); } diff --git a/src/org/apache/cassandra/db/SuperColumn.java b/src/org/apache/cassandra/db/SuperColumn.java index 34e597ed16..606e865686 100644 --- a/src/org/apache/cassandra/db/SuperColumn.java +++ b/src/org/apache/cassandra/db/SuperColumn.java @@ -49,6 +49,7 @@ public final class SuperColumn implements IColumn, Serializable private String name_; private EfficientBidiMap columns_ = new EfficientBidiMap(ColumnComparatorFactory.getComparator(ColumnComparatorFactory.ComparatorType.TIMESTAMP)); + private int localDeletionTime = Integer.MIN_VALUE; private long markedForDeleteAt = Long.MIN_VALUE; private AtomicInteger size_ = new AtomicInteger(0); @@ -289,7 +290,14 @@ public final class SuperColumn implements IColumn, Serializable return sb.toString(); } - public void markForDeleteAt(long timestamp) { + public int getLocalDeletionTime() + { + return localDeletionTime; + } + + public void markForDeleteAt(int localDeleteTime, long timestamp) + { + this.localDeletionTime = localDeleteTime; this.markedForDeleteAt = timestamp; } } @@ -300,20 +308,14 @@ class SuperColumnSerializer implements ICompactSerializer2 { SuperColumn superColumn = (SuperColumn)column; dos.writeUTF(superColumn.name()); + dos.writeInt(superColumn.getLocalDeletionTime()); dos.writeLong(superColumn.getMarkedForDeleteAt()); Collection columns = column.getSubColumns(); int size = columns.size(); dos.writeInt(size); - /* - * Add the total size of the columns. This is useful - * to skip over all the columns in this super column - * if we are not interested in this super column. - */ dos.writeInt(superColumn.getSizeOfAllColumns()); - // dos.writeInt(superColumn.size()); - for ( IColumn subColumn : columns ) { Column.serializer().serialize(subColumn, dos); @@ -328,7 +330,7 @@ class SuperColumnSerializer implements ICompactSerializer2 { String name = dis.readUTF(); SuperColumn superColumn = new SuperColumn(name); - superColumn.markForDeleteAt(dis.readLong()); + superColumn.markForDeleteAt(dis.readInt(), dis.readLong()); return superColumn; } diff --git a/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java b/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java index d6b6fc56fd..c3511c52f8 100644 --- a/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java +++ b/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java @@ -19,6 +19,8 @@ import org.apache.commons.lang.StringUtils; import org.apache.cassandra.ServerTest; import org.testng.annotations.Test; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNull; public class ColumnFamilyStoreTest extends ServerTest { @@ -207,7 +209,8 @@ public class ColumnFamilyStoreTest extends ServerTest rm.apply(); ColumnFamily retrieved = store.getColumnFamily("key1", "Standard1", new IdentityFilter()); - assert retrieved.getColumnCount() == 0; + assert retrieved.getColumn("Column1").isMarkedForDelete(); + assertNull(ColumnFamilyStore.removeDeleted(retrieved, Integer.MAX_VALUE)); } @Test @@ -229,7 +232,8 @@ public class ColumnFamilyStoreTest extends ServerTest rm.apply(); ColumnFamily retrieved = store.getColumnFamily("key1", "Super1:SC1", new IdentityFilter()); - assert retrieved.getColumnCount() == 0; + assert retrieved.getColumn("SC1").getSubColumn("Column1").isMarkedForDelete(); + assertNull(ColumnFamilyStore.removeDeleted(retrieved, Integer.MAX_VALUE)); } @Test @@ -258,7 +262,7 @@ public class ColumnFamilyStoreTest extends ServerTest Collection subColumns = resolved.getAllColumns().first().getSubColumns(); assert subColumns.size() == 1; assert subColumns.iterator().next().timestamp() == 0; - assert ColumnFamilyStore.removeDeleted(resolved).getColumnCount() == 0; + assertNull(ColumnFamilyStore.removeDeleted(resolved, Integer.MAX_VALUE)); } @Test @@ -281,7 +285,9 @@ public class ColumnFamilyStoreTest extends ServerTest rm.apply(); ColumnFamily retrieved = store.getColumnFamily("key1", "Standard1", new IdentityFilter()); - assert retrieved.getColumnCount() == 0; + assert retrieved.isMarkedForDelete(); + assertEquals(retrieved.getColumnCount(), 0); + assertNull(ColumnFamilyStore.removeDeleted(retrieved, Integer.MAX_VALUE)); } @Test