From 90170d1594c3a88f0a7b6a25da7161bb7af2e552 Mon Sep 17 00:00:00 2001 From: Pavel Yaskevich Date: Thu, 17 May 2012 22:35:22 +0300 Subject: [PATCH] Avoid ID conflicts from concurrent schema changes patch by Pavel Yaskevich; reviewed by Sylvain Lebresne for CASSANDRA-3794 --- CHANGES.txt | 1 + .../apache/cassandra/cache/RowCacheKey.java | 14 +- .../org/apache/cassandra/config/Avro.java | 6 +- .../apache/cassandra/config/CFMetaData.java | 46 ++++--- .../cassandra/config/DatabaseDescriptor.java | 1 - .../org/apache/cassandra/config/Schema.java | 66 ++++----- .../org/apache/cassandra/db/ColumnFamily.java | 10 +- .../cassandra/db/ColumnFamilySerializer.java | 62 +++++++-- .../cassandra/db/ColumnFamilyStore.java | 6 +- .../apache/cassandra/db/CounterMutation.java | 3 +- .../org/apache/cassandra/db/IMutation.java | 3 +- .../org/apache/cassandra/db/RowMutation.java | 38 ++--- src/java/org/apache/cassandra/db/Table.java | 22 +-- .../org/apache/cassandra/db/TypeSizes.java | 13 ++ .../db/UnknownColumnFamilyException.java | 5 +- .../cassandra/db/commitlog/CommitLog.java | 2 +- .../db/commitlog/CommitLogAllocator.java | 3 +- .../db/commitlog/CommitLogReplayer.java | 18 +-- .../db/commitlog/CommitLogSegment.java | 11 +- .../AbstractCompactionIterable.java | 5 +- .../db/compaction/CompactionInfo.java | 7 +- .../db/compaction/CompactionManager.java | 6 +- .../db/index/SecondaryIndexBuilder.java | 3 +- .../cassandra/streaming/StreamRequest.java | 7 +- .../apache/cassandra/utils/FBUtilities.java | 2 +- test/data/serialization/1.2/db.Row.bin | Bin 495 -> 527 bytes .../data/serialization/1.2/db.RowMutation.bin | Bin 3266 -> 3602 bytes .../1.2/streaming.StreamRequestMessage.bin | Bin 7127 -> 7151 bytes .../org/apache/cassandra/SchemaLoader.java | 130 +++++++++++------- test/unit/org/apache/cassandra/Util.java | 2 +- .../cassandra/cache/CacheProviderTest.java | 10 +- .../org/apache/cassandra/config/DefsTest.java | 19 +-- .../apache/cassandra/db/CommitLogTest.java | 7 +- .../cassandra/db/SerializationsTest.java | 15 +- 34 files changed, 312 insertions(+), 231 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 20a4e1fdb1..82ce164132 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -8,6 +8,7 @@ * Save IndexSummary into new SSTable 'Summary' component (CASSANDRA-2392) * Add support for range tombstones (CASSANDRA-3708) * Improve MessagingService efficiency (CASSANDRA-3617) + * Avoid ID conflicts from concurrent schema changes (CASSANDRA-3794) 1.1.1-dev diff --git a/src/java/org/apache/cassandra/cache/RowCacheKey.java b/src/java/org/apache/cassandra/cache/RowCacheKey.java index 44d1c64544..009483c6d4 100644 --- a/src/java/org/apache/cassandra/cache/RowCacheKey.java +++ b/src/java/org/apache/cassandra/cache/RowCacheKey.java @@ -21,6 +21,7 @@ import java.io.DataOutputStream; import java.io.IOException; import java.nio.ByteBuffer; import java.util.Arrays; +import java.util.UUID; import org.apache.cassandra.config.Schema; import org.apache.cassandra.db.TypeSizes; @@ -31,15 +32,15 @@ import org.apache.cassandra.utils.Pair; public class RowCacheKey implements CacheKey, Comparable { - public final int cfId; + public final UUID cfId; public final byte[] key; - public RowCacheKey(int cfId, DecoratedKey key) + public RowCacheKey(UUID cfId, DecoratedKey key) { this(cfId, key.key); } - public RowCacheKey(int cfId, ByteBuffer key) + public RowCacheKey(UUID cfId, ByteBuffer key) { this.cfId = cfId; this.key = ByteBufferUtil.getArray(key); @@ -69,21 +70,20 @@ public class RowCacheKey implements CacheKey, Comparable RowCacheKey that = (RowCacheKey) o; - if (cfId != that.cfId) return false; - return Arrays.equals(key, that.key); + return cfId.equals(that.cfId) && Arrays.equals(key, that.key); } @Override public int hashCode() { - int result = cfId; + int result = cfId.hashCode(); result = 31 * result + (key != null ? Arrays.hashCode(key) : 0); return result; } public int compareTo(RowCacheKey otherKey) { - return (cfId < otherKey.cfId) ? -1 : ((cfId == otherKey.cfId) ? FBUtilities.compareUnsigned(key, otherKey.key, 0, 0, key.length, otherKey.key.length) : 1); + return (cfId.compareTo(otherKey.cfId) < 0) ? -1 : ((cfId.equals(otherKey.cfId)) ? FBUtilities.compareUnsigned(key, otherKey.key, 0, 0, key.length, otherKey.key.length) : 1); } @Override diff --git a/src/java/org/apache/cassandra/config/Avro.java b/src/java/org/apache/cassandra/config/Avro.java index 799634543f..9991cd0f9f 100644 --- a/src/java/org/apache/cassandra/config/Avro.java +++ b/src/java/org/apache/cassandra/config/Avro.java @@ -127,8 +127,7 @@ public class Avro cf.name.toString(), ColumnFamilyType.create(cf.column_type.toString()), comparator, - subcolumnComparator, - cf.id); + subcolumnComparator); // When we pull up an old avro CfDef which doesn't have these arguments, // it doesn't default them correctly. Without explicit defaulting, @@ -178,6 +177,9 @@ public class Avro throw new RuntimeException(e); } + // adding old -> new style ID mapping to support backward compatibility + Schema.instance.addOldCfIdMapping(cf.id, newCFMD.cfId); + return newCFMD.comment(cf.comment.toString()) .readRepairChance(cf.read_repair_chance) .dcLocalReadRepairChance(cf.dclocal_read_repair_chance) diff --git a/src/java/org/apache/cassandra/config/CFMetaData.java b/src/java/org/apache/cassandra/config/CFMetaData.java index cac2931de9..d128719e42 100644 --- a/src/java/org/apache/cassandra/config/CFMetaData.java +++ b/src/java/org/apache/cassandra/config/CFMetaData.java @@ -21,10 +21,12 @@ import java.io.IOException; import java.lang.reflect.Constructor; import java.lang.reflect.InvocationTargetException; import java.nio.ByteBuffer; +import java.security.MessageDigest; import java.util.*; import com.google.common.collect.MapDifference; import com.google.common.collect.Maps; +import org.apache.commons.lang.ArrayUtils; import org.apache.commons.lang.StringUtils; import org.apache.commons.lang.builder.EqualsBuilder; import org.apache.commons.lang.builder.HashCodeBuilder; @@ -179,7 +181,7 @@ public final class CFMetaData } //REQUIRED - public final Integer cfId; // internal id, never exposed to user + public final UUID cfId; // internal id, never exposed to user public final String ksName; // name of keyspace public final String cfName; // name of this column family public final ColumnFamilyType cfType; // standard, super @@ -243,10 +245,10 @@ public final class CFMetaData public CFMetaData(String keyspace, String name, ColumnFamilyType type, AbstractType comp, AbstractType subcc) { - this(keyspace, name, type, comp, subcc, Schema.instance.nextCFId()); + this(keyspace, name, type, comp, subcc, getId(keyspace, name)); } - CFMetaData(String keyspace, String name, ColumnFamilyType type, AbstractType comp, AbstractType subcc, int id) + CFMetaData(String keyspace, String name, ColumnFamilyType type, AbstractType comp, AbstractType subcc, UUID id) { // Final fields must be set in constructor ksName = keyspace; @@ -254,9 +256,6 @@ public final class CFMetaData cfType = type; comparator = comp; subcolumnComparator = enforceSubccDefault(type, subcc); - - // System cfs have specific ids, and copies of old CFMDs need - // to copy over the old id. cfId = id; this.init(); @@ -272,6 +271,11 @@ public final class CFMetaData return (comment == null) ? "" : comment.toString(); } + static UUID getId(String ksName, String cfName) + { + return UUID.nameUUIDFromBytes(ArrayUtils.addAll(ksName.getBytes(), cfName.getBytes())); + } + private void init() { // Set a bunch of defaults @@ -306,10 +310,13 @@ public final class CFMetaData updateCfDef(); // init cqlCfDef } - private static CFMetaData newSystemMetadata(String cfName, int cfId, String comment, AbstractType comparator, AbstractType subcc) + private static CFMetaData newSystemMetadata(String cfName, int oldCfId, String comment, AbstractType comparator, AbstractType subcc) { ColumnFamilyType type = subcc == null ? ColumnFamilyType.Standard : ColumnFamilyType.Super; - CFMetaData newCFMD = new CFMetaData(Table.SYSTEM_TABLE, cfName, type, comparator, subcc, cfId); + CFMetaData newCFMD = new CFMetaData(Table.SYSTEM_TABLE, cfName, type, comparator, subcc); + + // adding old -> new style ID mapping to support backward compatibility + Schema.instance.addOldCfIdMapping(oldCfId, newCFMD.cfId); return newCFMD.comment(comment) .readRepairChance(0) @@ -317,7 +324,7 @@ public final class CFMetaData .gcGraceSeconds(0); } - private static CFMetaData newSchemaMetadata(String cfName, int cfId, String comment, AbstractType comparator, AbstractType subcc) + private static CFMetaData newSchemaMetadata(String cfName, int oldCfId, String comment, AbstractType comparator, AbstractType subcc) { /* * Schema column families needs a gc_grace (since they are replicated @@ -325,7 +332,7 @@ public final class CFMetaData * could be dead for that long a time. */ int gcGrace = 120 * 24 * 3600; // 3 months - return newSystemMetadata(cfName, cfId, comment, comparator, subcc).gcGraceSeconds(gcGrace); + return newSystemMetadata(cfName, oldCfId, comment, comparator, subcc).gcGraceSeconds(gcGrace); } public static CFMetaData newIndexMetadata(CFMetaData parent, ColumnDefinition info, AbstractType columnComparator) @@ -527,7 +534,7 @@ public final class CFMetaData .append(keyValidator, rhs.keyValidator) .append(minCompactionThreshold, rhs.minCompactionThreshold) .append(maxCompactionThreshold, rhs.maxCompactionThreshold) - .append(cfId.intValue(), rhs.cfId.intValue()) + .append(cfId, rhs.cfId) .append(column_metadata, rhs.column_metadata) .append(keyAlias, rhs.keyAlias) .append(columnAliases, rhs.columnAliases) @@ -623,8 +630,7 @@ public final class CFMetaData cf_def.name, cfType, TypeParser.parse(cf_def.comparator_type), - cf_def.subcomparator_type == null ? null : TypeParser.parse(cf_def.subcomparator_type), - cf_def.isSetId() ? cf_def.id : Schema.instance.nextCFId()); + cf_def.subcomparator_type == null ? null : TypeParser.parse(cf_def.subcomparator_type)); if (cf_def.isSetGc_grace_seconds()) { newCFMD.gcGraceSeconds(cf_def.gc_grace_seconds); } if (cf_def.isSetMin_compaction_threshold()) { newCFMD.minCompactionThreshold(cf_def.min_compaction_threshold); } @@ -799,7 +805,6 @@ public final class CFMetaData public org.apache.cassandra.thrift.CfDef toThrift() { org.apache.cassandra.thrift.CfDef def = new org.apache.cassandra.thrift.CfDef(ksName, cfName); - def.setId(cfId); def.setColumn_type(cfType.name()); def.setComparator_type(comparator.toString()); if (subcolumnComparator != null) @@ -1142,7 +1147,11 @@ public final class CFMetaData ColumnFamily cf = rm.addOrGet(SystemTable.SCHEMA_COLUMNFAMILIES_CF); int ldt = (int) (System.currentTimeMillis() / 1000); - cf.addColumn(Column.create(cfId, timestamp, cfName, "id")); + Integer oldId = Schema.instance.convertNewCfId(cfId); + + if (oldId != null) // keep old ids (see CASSANDRA-3794 for details) + cf.addColumn(Column.create(oldId, timestamp, cfName, "id")); + cf.addColumn(Column.create(cfType.toString(), timestamp, cfName, "type")); cf.addColumn(Column.create(comparator.toString(), timestamp, cfName, "comparator")); if (subcolumnComparator != null) @@ -1179,8 +1188,11 @@ public final class CFMetaData result.getString("columnfamily"), ColumnFamilyType.valueOf(result.getString("type")), TypeParser.parse(result.getString("comparator")), - result.has("subcomparator") ? TypeParser.parse(result.getString("subcomparator")) : null, - result.getInt("id")); + result.has("subcomparator") ? TypeParser.parse(result.getString("subcomparator")) : null); + + if (result.has("id"))// try to identify if ColumnFamily Id is old style (before C* 1.2) and add old -> new mapping if so + Schema.instance.addOldCfIdMapping(result.getInt("id"), cfm.cfId); + cfm.readRepairChance(result.getDouble("read_repair_chance")); cfm.dcLocalReadRepairChance(result.getDouble("local_read_repair_chance")); cfm.replicateOnWrite(result.getBoolean("replicate_on_write")); diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index ce6bbe76f0..6d85425495 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -536,7 +536,6 @@ public class DatabaseDescriptor } Schema.instance.updateVersion(); - Schema.instance.fixCFMaxId(); } private static boolean hasExistingNoSystemTables() diff --git a/src/java/org/apache/cassandra/config/Schema.java b/src/java/org/apache/cassandra/config/Schema.java index 32ec07addb..7d52d2e63d 100644 --- a/src/java/org/apache/cassandra/config/Schema.java +++ b/src/java/org/apache/cassandra/config/Schema.java @@ -17,7 +17,6 @@ */ package org.apache.cassandra.config; -import java.io.IOError; import java.nio.ByteBuffer; import java.security.MessageDigest; import java.util.*; @@ -55,9 +54,6 @@ public class Schema */ public static final int NAME_LENGTH = 48; - private static final int MIN_CF_ID = 1000; - private final AtomicInteger cfIdGen = new AtomicInteger(MIN_CF_ID); - /* metadata map for faster table lookup */ private final Map tables = new NonBlockingHashMap(); @@ -65,7 +61,9 @@ public class Schema private final Map tableInstances = new NonBlockingHashMap(); /* metadata map for faster ColumnFamily lookup */ - private final BiMap, Integer> cfIdMap = HashBiMap.create(); + private final BiMap, UUID> cfIdMap = HashBiMap.create(); + // mapping from old ColumnFamily Id (Integer) to a new version which is UUID + private final BiMap oldCfIdMap = HashBiMap.create(); private volatile UUID version; private final ReadWriteLock versionLock = new ReentrantReadWriteLock(); @@ -106,8 +104,6 @@ public class Schema setTableDefinition(keyspaceDef); - fixCFMaxId(); - return this; } @@ -184,7 +180,7 @@ public class Schema * * @return metadata about ColumnFamily */ - public CFMetaData getCFMetaData(Integer cfId) + public CFMetaData getCFMetaData(UUID cfId) { Pair cf = getCF(cfId); return (cf == null) ? null : getCFMetaData(cf.left, cf.right); @@ -343,11 +339,34 @@ public class Schema /* ColumnFamily query/control methods */ + public void addOldCfIdMapping(Integer oldId, UUID newId) + { + if (oldId == null) + return; + + oldCfIdMap.put(oldId, newId); + } + + public UUID convertOldCfId(Integer oldCfId) + { + UUID cfId = oldCfIdMap.get(oldCfId); + + if (cfId == null) + throw new IllegalArgumentException("ColumnFamily identified by old " + oldCfId + " was not found."); + + return cfId; + } + + public Integer convertNewCfId(UUID newCfId) + { + return oldCfIdMap.containsValue(newCfId) ? oldCfIdMap.inverse().get(newCfId) : null; + } + /** * @param cfId The identifier of the ColumnFamily to lookup * @return The (ksname,cfname) pair for the given id, or null if it has been dropped. */ - public Pair getCF(Integer cfId) + public Pair getCF(UUID cfId) { return cfIdMap.inverse().get(cfId); } @@ -360,7 +379,7 @@ public class Schema * * @return The id for the given (ksname,cfname) pair, or null if it has been dropped. */ - public Integer getId(String ksName, String cfName) + public UUID getId(String ksName, String cfName) { return cfIdMap.get(new Pair(ksName, cfName)); } @@ -382,8 +401,6 @@ public class Schema logger.debug("Adding {} to cfIdMap", cfm); cfIdMap.put(key, cfm.cfId); - - fixCFMaxId(); } /** @@ -396,30 +413,6 @@ public class Schema cfIdMap.remove(new Pair(cfm.ksName, cfm.cfName)); } - /** - * This gets called after initialization to make sure that id generation happens properly. - */ - public void fixCFMaxId() - { - int cval, nval; - do - { - cval = cfIdGen.get(); - int inMap = cfIdMap.isEmpty() ? 0 : Collections.max(cfIdMap.values()) + 1; - // never set it to less than 1000. this ensures that we have enough system CFids for future use. - nval = Math.max(Math.max(inMap, cval), MIN_CF_ID); - } - while (!cfIdGen.compareAndSet(cval, nval)); - } - - /** - * @return identifier for the new ColumnFamily (called primarily by CFMetaData constructor) - */ - public int nextCFId() - { - return cfIdGen.getAndIncrement(); - } - /* Version control */ /** @@ -494,6 +487,5 @@ public class Schema } updateVersionAndAnnounce(); - fixCFMaxId(); } } diff --git a/src/java/org/apache/cassandra/db/ColumnFamily.java b/src/java/org/apache/cassandra/db/ColumnFamily.java index 71378eb519..c2dd118f0c 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamily.java +++ b/src/java/org/apache/cassandra/db/ColumnFamily.java @@ -19,9 +19,11 @@ package org.apache.cassandra.db; import java.nio.ByteBuffer; import java.security.MessageDigest; +import java.util.UUID; import org.apache.cassandra.io.sstable.SSTable; import org.apache.cassandra.utils.*; + import org.apache.commons.lang.builder.HashCodeBuilder; import org.apache.cassandra.cache.IRowCacheEntry; @@ -32,8 +34,6 @@ import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.db.marshal.MarshalException; import org.apache.cassandra.io.IColumnSerializer; import org.apache.cassandra.io.sstable.ColumnStats; -import org.apache.cassandra.io.sstable.SSTable; -import org.apache.cassandra.utils.*; public class ColumnFamily extends AbstractColumnContainer implements IRowCacheEntry { @@ -41,12 +41,12 @@ public class ColumnFamily extends AbstractColumnContainer implements IRowCacheEn public static final ColumnFamilySerializer serializer = new ColumnFamilySerializer(); private final CFMetaData cfm; - public static ColumnFamily create(Integer cfId) + public static ColumnFamily create(UUID cfId) { return create(Schema.instance.getCFMetaData(cfId)); } - public static ColumnFamily create(Integer cfId, ISortedColumns.Factory factory) + public static ColumnFamily create(UUID cfId, ISortedColumns.Factory factory) { return create(Schema.instance.getCFMetaData(cfId), factory); } @@ -108,7 +108,7 @@ public class ColumnFamily extends AbstractColumnContainer implements IRowCacheEn return cf; } - public Integer id() + public UUID id() { return cfm.cfId; } diff --git a/src/java/org/apache/cassandra/db/ColumnFamilySerializer.java b/src/java/org/apache/cassandra/db/ColumnFamilySerializer.java index cfc4b637be..831fd3db12 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilySerializer.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilySerializer.java @@ -20,6 +20,8 @@ package org.apache.cassandra.db; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; +import java.util.UUID; + import org.apache.cassandra.config.Schema; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -28,6 +30,8 @@ import org.apache.cassandra.io.IColumnSerializer; import org.apache.cassandra.io.IVersionedSerializer; import org.apache.cassandra.io.ISSTableSerializer; import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.net.MessagingService; +import org.apache.cassandra.utils.UUIDGen; public class ColumnFamilySerializer implements IVersionedSerializer, ISSTableSerializer { @@ -61,7 +65,7 @@ public class ColumnFamilySerializer implements IVersionedSerializer uuid mapping could not be established (CF was created in mixed version cluster)."); + + dos.writeInt(oldId); + } + else + UUIDGen.serializer.serialize(cfId, dos, version); + } + + public UUID deserializeCfId(DataInput dis, int version) throws IOException + { + // create a ColumnFamily based on the cf id + UUID cfId = (version < MessagingService.VERSION_12) + ? Schema.instance.convertOldCfId(dis.readInt()) + : UUIDGen.serializer.deserialize(dis, version); + + if (Schema.instance.getCF(cfId) == null) + throw new UnknownColumnFamilyException("Couldn't find cfId=" + cfId, cfId); + + return cfId; + } + + public int cfIdSerializedSize(UUID cfId, TypeSizes typeSizes, int version) + { + if (version < MessagingService.VERSION_12) // try to use CF's old id where possible (CASSANDRA-3794) + { + Integer oldId = Schema.instance.convertNewCfId(cfId); + + if (oldId == null) + throw new RuntimeException("Can't serialize ColumnFamily ID " + cfId + " to be used by version " + version + + ", because int <-> uuid mapping could not be established (CF was created in mixed version cluster)."); + + return typeSizes.sizeof(oldId); + } + + return typeSizes.sizeof(cfId); + } } diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 2f89b0464f..2891585477 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -1122,7 +1122,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean * @return the entire row for filter.key, if present in the cache (or we can cache it), or just the column * specified by filter otherwise */ - private ColumnFamily getThroughCache(Integer cfId, QueryFilter filter) + private ColumnFamily getThroughCache(UUID cfId, QueryFilter filter) { assert isRowCacheEnabled() : String.format("Row cache is not enabled on column family [" + getColumnFamilyName() + "]"); @@ -1181,7 +1181,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return cf.isSuper() ? removeDeleted(cf, gcBefore) : removeDeletedCF(cf, gcBefore); } - Integer cfId = Schema.instance.getId(table.name, this.columnFamily); + UUID cfId = Schema.instance.getId(table.name, this.columnFamily); if (cfId == null) return null; // secondary index @@ -1586,7 +1586,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean public void invalidateCachedRow(DecoratedKey key) { - Integer cfId = Schema.instance.getId(table.name, this.columnFamily); + UUID cfId = Schema.instance.getId(table.name, this.columnFamily); if (cfId == null) return; // secondary index diff --git a/src/java/org/apache/cassandra/db/CounterMutation.java b/src/java/org/apache/cassandra/db/CounterMutation.java index c3256cc24d..07ae713b7b 100644 --- a/src/java/org/apache/cassandra/db/CounterMutation.java +++ b/src/java/org/apache/cassandra/db/CounterMutation.java @@ -24,6 +24,7 @@ import java.nio.ByteBuffer; import java.util.Collection; import java.util.LinkedList; import java.util.List; +import java.util.UUID; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -55,7 +56,7 @@ public class CounterMutation implements IMutation return rowMutation.getTable(); } - public Collection getColumnFamilyIds() + public Collection getColumnFamilyIds() { return rowMutation.getColumnFamilyIds(); } diff --git a/src/java/org/apache/cassandra/db/IMutation.java b/src/java/org/apache/cassandra/db/IMutation.java index 119dfb57ff..cd08b46d20 100644 --- a/src/java/org/apache/cassandra/db/IMutation.java +++ b/src/java/org/apache/cassandra/db/IMutation.java @@ -20,11 +20,12 @@ package org.apache.cassandra.db; import java.nio.ByteBuffer; import java.io.IOException; import java.util.Collection; +import java.util.UUID; public interface IMutation { public String getTable(); - public Collection getColumnFamilyIds(); + public Collection getColumnFamilyIds(); public ByteBuffer key(); public void apply() throws IOException; public String toString(boolean shallow); diff --git a/src/java/org/apache/cassandra/db/RowMutation.java b/src/java/org/apache/cassandra/db/RowMutation.java index 5f3c9d1cb9..0d83cec100 100644 --- a/src/java/org/apache/cassandra/db/RowMutation.java +++ b/src/java/org/apache/cassandra/db/RowMutation.java @@ -50,13 +50,13 @@ public class RowMutation implements IMutation private final String table; private final ByteBuffer key; // map of column family id to mutations for that column family. - protected final Map modifications; + protected Map modifications = new HashMap(); private final Map preserializedBuffers = new HashMap(); public RowMutation(String table, ByteBuffer key) { - this(table, key, new HashMap()); + this(table, key, new HashMap()); } public RowMutation(String table, Row row) @@ -65,7 +65,7 @@ public class RowMutation implements IMutation add(row.cf); } - protected RowMutation(String table, ByteBuffer key, Map modifications) + protected RowMutation(String table, ByteBuffer key, Map modifications) { this.table = table; this.key = key; @@ -77,7 +77,7 @@ public class RowMutation implements IMutation return table; } - public Collection getColumnFamilyIds() + public Collection getColumnFamilyIds() { return modifications.keySet(); } @@ -92,7 +92,7 @@ public class RowMutation implements IMutation return modifications.values(); } - public ColumnFamily getColumnFamily(Integer cfId) + public ColumnFamily getColumnFamily(UUID cfId) { return modifications.get(cfId); } @@ -199,8 +199,9 @@ public class RowMutation implements IMutation */ public void add(QueryPath path, ByteBuffer value, long timestamp, int timeToLive) { - Integer id = Schema.instance.getId(table, path.columnFamilyName); + UUID id = Schema.instance.getId(table, path.columnFamilyName); ColumnFamily columnFamily = modifications.get(id); + if (columnFamily == null) { columnFamily = ColumnFamily.create(table, path.columnFamilyName); @@ -211,8 +212,9 @@ public class RowMutation implements IMutation public void addCounter(QueryPath path, long value) { - Integer id = Schema.instance.getId(table, path.columnFamilyName); + UUID id = Schema.instance.getId(table, path.columnFamilyName); ColumnFamily columnFamily = modifications.get(id); + if (columnFamily == null) { columnFamily = ColumnFamily.create(table, path.columnFamilyName); @@ -228,7 +230,7 @@ public class RowMutation implements IMutation public void delete(QueryPath path, long timestamp) { - Integer id = Schema.instance.getId(table, path.columnFamilyName); + UUID id = Schema.instance.getId(table, path.columnFamilyName); int localDeleteTime = (int) (System.currentTimeMillis() / 1000); @@ -264,7 +266,7 @@ public class RowMutation implements IMutation if (!table.equals(rm.table) || !key.equals(rm.key)) throw new IllegalArgumentException(); - for (Map.Entry entry : rm.modifications.entrySet()) + for (Map.Entry entry : rm.modifications.entrySet()) { // It's slighty faster to assume the key wasn't present and fix if // not in the case where it wasn't there indeed. @@ -325,7 +327,7 @@ public class RowMutation implements IMutation if (shallow) { List cfnames = new ArrayList(modifications.size()); - for (Integer cfid : modifications.keySet()) + for (UUID cfid : modifications.keySet()) { CFMetaData cfm = Schema.instance.getCFMetaData(cfid); cfnames.add(cfm == null ? "-dropped-" : cfm.cfName); @@ -385,7 +387,7 @@ public class RowMutation implements IMutation { RowMutation rm = serializer.deserialize(new DataInputStream(new FastByteArrayInputStream(raw)), version); boolean hasCounters = false; - for (Map.Entry entry : rm.modifications.entrySet()) + for (Map.Entry entry : rm.modifications.entrySet()) { if (entry.getValue().metadata().getDefaultValidator().isCommutative()) { @@ -411,9 +413,9 @@ public class RowMutation implements IMutation int size = rm.modifications.size(); dos.writeInt(size); assert size >= 0; - for (Map.Entry entry : rm.modifications.entrySet()) + for (Map.Entry entry : rm.modifications.entrySet()) { - dos.writeInt(entry.getKey()); + ColumnFamily.serializer.serializeCfId(entry.getKey(), dos, version); ColumnFamily.serializer.serialize(entry.getValue(), dos, version); } } @@ -422,13 +424,13 @@ public class RowMutation implements IMutation { String table = dis.readUTF(); ByteBuffer key = ByteBufferUtil.readWithShortLength(dis); - Map modifications = new HashMap(); + Map modifications = new HashMap(); int size = dis.readInt(); for (int i = 0; i < size; ++i) { - Integer cfid = Integer.valueOf(dis.readInt()); + UUID cfId = ColumnFamily.serializer.deserializeCfId(dis, version); ColumnFamily cf = ColumnFamily.serializer.deserialize(dis, flag, TreeMapBackedSortedColumns.factory(), version); - modifications.put(cfid, cf); + modifications.put(cfId, cf); } return new RowMutation(table, key, modifications); } @@ -446,9 +448,9 @@ public class RowMutation implements IMutation size += sizes.sizeof((short) keySize) + keySize; size += sizes.sizeof(rm.modifications.size()); - for (Map.Entry entry : rm.modifications.entrySet()) + for (Map.Entry entry : rm.modifications.entrySet()) { - size += sizes.sizeof(entry.getKey()); + size += ColumnFamily.serializer.cfIdSerializedSize(entry.getValue().id(), sizes, version); size += ColumnFamily.serializer.serializedSize(entry.getValue(), TypeSizes.NATIVE, version); } diff --git a/src/java/org/apache/cassandra/db/Table.java b/src/java/org/apache/cassandra/db/Table.java index e4a26c184b..901ba34435 100644 --- a/src/java/org/apache/cassandra/db/Table.java +++ b/src/java/org/apache/cassandra/db/Table.java @@ -20,15 +20,7 @@ package org.apache.cassandra.db; import java.io.IOError; import java.io.IOException; import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.Collection; -import java.util.Collections; -import java.util.Comparator; -import java.util.Iterator; -import java.util.List; -import java.util.Map; -import java.util.SortedSet; -import java.util.TreeSet; +import java.util.*; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; @@ -87,7 +79,7 @@ public class Table /* Table name. */ public final String name; /* ColumnFamilyStore per column family */ - private final Map columnFamilyStores = new ConcurrentHashMap(); + private final Map columnFamilyStores = new ConcurrentHashMap(); private final Object[] indexLocks; private volatile AbstractReplicationStrategy replicationStrategy; @@ -148,13 +140,13 @@ public class Table public ColumnFamilyStore getColumnFamilyStore(String cfName) { - Integer id = Schema.instance.getId(name, cfName); + UUID id = Schema.instance.getId(name, cfName); if (id == null) throw new IllegalArgumentException(String.format("Unknown table/cf pair (%s.%s)", name, cfName)); return getColumnFamilyStore(id); } - public ColumnFamilyStore getColumnFamilyStore(Integer id) + public ColumnFamilyStore getColumnFamilyStore(UUID id) { ColumnFamilyStore cfs = columnFamilyStores.get(id); if (cfs == null) @@ -313,7 +305,7 @@ public class Table } // best invoked on the compaction mananger. - public void dropCf(Integer cfId) throws IOException + public void dropCf(UUID cfId) throws IOException { assert columnFamilyStores.containsKey(cfId); ColumnFamilyStore cfs = columnFamilyStores.remove(cfId); @@ -342,7 +334,7 @@ public class Table } /** adds a cf to internal structures, ends up creating disk files). */ - public void initCf(Integer cfId, String cfName) + public void initCf(UUID cfId, String cfName) { if (columnFamilyStores.containsKey(cfId)) { @@ -556,7 +548,7 @@ public class Table public List> flush() throws IOException { List> futures = new ArrayList>(); - for (Integer cfId : columnFamilyStores.keySet()) + for (UUID cfId : columnFamilyStores.keySet()) { Future future = columnFamilyStores.get(cfId).forceFlush(); if (future != null) diff --git a/src/java/org/apache/cassandra/db/TypeSizes.java b/src/java/org/apache/cassandra/db/TypeSizes.java index aac89d0ddf..67f8fccf6d 100644 --- a/src/java/org/apache/cassandra/db/TypeSizes.java +++ b/src/java/org/apache/cassandra/db/TypeSizes.java @@ -18,6 +18,7 @@ package org.apache.cassandra.db; import java.nio.ByteBuffer; +import java.util.UUID; import org.apache.cassandra.utils.FBUtilities; @@ -30,11 +31,13 @@ public abstract class TypeSizes private static final int SHORT_SIZE = 2; private static final int INT_SIZE = 4; private static final int LONG_SIZE = 8; + private static final int UUID_SIZE = 16; public abstract int sizeof(boolean value); public abstract int sizeof(short value); public abstract int sizeof(int value); public abstract int sizeof(long value); + public abstract int sizeof(UUID value); /** assumes UTF8 */ public int sizeof(String value) @@ -92,6 +95,11 @@ public abstract class TypeSizes { return LONG_SIZE; } + + public int sizeof(UUID value) + { + return UUID_SIZE; + } } public static class VIntEncodedTypeSizes extends TypeSizes @@ -141,5 +149,10 @@ public abstract class TypeSizes { return sizeofVInt(i); } + + public int sizeof(UUID value) + { + return sizeofVInt(value.getMostSignificantBits()) + sizeofVInt(value.getLeastSignificantBits()); + } } } diff --git a/src/java/org/apache/cassandra/db/UnknownColumnFamilyException.java b/src/java/org/apache/cassandra/db/UnknownColumnFamilyException.java index e1b2983575..c43b50a6b9 100644 --- a/src/java/org/apache/cassandra/db/UnknownColumnFamilyException.java +++ b/src/java/org/apache/cassandra/db/UnknownColumnFamilyException.java @@ -18,13 +18,14 @@ package org.apache.cassandra.db; import java.io.IOException; +import java.util.UUID; public class UnknownColumnFamilyException extends IOException { - public final int cfId; + public final UUID cfId; - public UnknownColumnFamilyException(String msg, int cfId) + public UnknownColumnFamilyException(String msg, UUID cfId) { super(msg); this.cfId = cfId; diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLog.java b/src/java/org/apache/cassandra/db/commitlog/CommitLog.java index 055a32d347..080cd2260b 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLog.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLog.java @@ -208,7 +208,7 @@ public class CommitLog implements CommitLogMBean * @param cfId the column family ID that was flushed * @param context the replay position of the flush */ - public void discardCompletedSegments(final Integer cfId, final ReplayPosition context) throws IOException + public void discardCompletedSegments(final UUID cfId, final ReplayPosition context) throws IOException { Callable task = new Callable() { diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogAllocator.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogAllocator.java index 4a3bf0a9e9..e402a4a79f 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogAllocator.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogAllocator.java @@ -23,6 +23,7 @@ import java.io.IOError; import java.io.IOException; import java.util.Collection; import java.util.Collections; +import java.util.UUID; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; @@ -300,7 +301,7 @@ public class CommitLogAllocator if (oldestSegment != null) { - for (Integer dirtyCFId : oldestSegment.getDirtyCFIDs()) + for (UUID dirtyCFId : oldestSegment.getDirtyCFIDs()) { String keypace = Schema.instance.getCF(dirtyCFId).left; final ColumnFamilyStore cfs = Table.open(keypace).getColumnFamilyStore(dirtyCFId); diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java index 812c774614..69853597c1 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogReplayer.java @@ -4,11 +4,7 @@ import java.io.DataInputStream; import java.io.EOFException; import java.io.File; import java.io.IOException; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Set; +import java.util.*; import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicInteger; import java.util.zip.Checksum; @@ -44,9 +40,9 @@ public class CommitLogReplayer private final Set tablesRecovered; private final List> futures; - private final Map invalidMutations; + private final Map invalidMutations; private final AtomicInteger replayedCount; - private final Map cfPositions; + private final Map cfPositions; private final ReplayPosition globalPosition; private final Checksum checksum; private byte[] buffer; @@ -56,11 +52,11 @@ private final AtomicInteger replayedCount; this.tablesRecovered = new NonBlockingHashSet
(); this.futures = new ArrayList>(); this.buffer = new byte[4096]; - this.invalidMutations = new HashMap(); + this.invalidMutations = new HashMap(); // count the number of replayed mutation. We don't really care about atomicity, but we need it to be a reference. this.replayedCount = new AtomicInteger(); // compute per-CF and global replay positions - this.cfPositions = new HashMap(); + this.cfPositions = new HashMap(); for (ColumnFamilyStore cfs : ColumnFamilyStore.all()) { // it's important to call RP.gRP per-cf, before aggregating all the positions w/ the Ordering.min call @@ -81,8 +77,8 @@ private final AtomicInteger replayedCount; public int blockForWrites() throws IOException { - for (Map.Entry entry : invalidMutations.entrySet()) - logger.info(String.format("Skipped %d mutations from unknown (probably removed) CF with id %d", entry.getValue().intValue(), entry.getKey())); + for (Map.Entry entry : invalidMutations.entrySet()) + logger.info(String.format("Skipped %d mutations from unknown (probably removed) CF with id %s", entry.getValue().intValue(), entry.getKey())); // wait for all the writes to finish on the mutation stage FBUtilities.waitOnFutures(futures); diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java index be7c1ed9f8..0f3cb067d5 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java @@ -24,6 +24,7 @@ import java.io.RandomAccessFile; import java.nio.channels.FileChannel; import java.nio.MappedByteBuffer; import java.util.Collection; +import java.util.UUID; import java.util.regex.Matcher; import java.util.regex.Pattern; import java.util.zip.Checksum; @@ -58,7 +59,7 @@ public class CommitLogSegment static final int ENTRY_OVERHEAD_SIZE = 4 + 8 + 8; // cache which cf is dirty in this segment to avoid having to lookup all ReplayPositions to decide if we can delete this segment - private final HashMap cfLastWrite = new HashMap(); + private final HashMap cfLastWrite = new HashMap(); public final long id; @@ -326,7 +327,7 @@ public class CommitLogSegment * @param cfId the column family ID that is now dirty * @param position the position the last write for this CF was written at */ - private void markCFDirty(Integer cfId, Integer position) + private void markCFDirty(UUID cfId, Integer position) { cfLastWrite.put(cfId, position); } @@ -339,7 +340,7 @@ public class CommitLogSegment * @param cfId the column family ID that is now clean * @param context the optional clean offset */ - public void markClean(Integer cfId, ReplayPosition context) + public void markClean(UUID cfId, ReplayPosition context) { Integer lastWritten = cfLastWrite.get(cfId); @@ -352,7 +353,7 @@ public class CommitLogSegment /** * @return a collection of dirty CFIDs for this segment file. */ - public Collection getDirtyCFIDs() + public Collection getDirtyCFIDs() { return cfLastWrite.keySet(); } @@ -380,7 +381,7 @@ public class CommitLogSegment public String dirtyString() { StringBuilder sb = new StringBuilder(); - for (Integer cfId : cfLastWrite.keySet()) + for (UUID cfId : cfLastWrite.keySet()) { CFMetaData m = Schema.instance.getCFMetaData(cfId); sb.append(m == null ? "" : m.cfName).append(" (").append(cfId).append("), "); diff --git a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java index 41441dfed2..40f57a52a0 100644 --- a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java +++ b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionIterable.java @@ -49,10 +49,7 @@ public abstract class AbstractCompactionIterable extends CompactionInfo.Holder i public CompactionInfo getCompactionInfo() { - return new CompactionInfo(this.hashCode(), - type, - bytesRead, - totalBytes); + return new CompactionInfo(type, bytesRead, totalBytes); } public abstract CloseableIterator iterator(); diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java b/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java index 02b243360d..cc46c31204 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java @@ -20,6 +20,7 @@ package org.apache.cassandra.db.compaction; import java.io.Serializable; import java.util.HashMap; import java.util.Map; +import java.util.UUID; import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.Schema; @@ -39,7 +40,7 @@ public final class CompactionInfo implements Serializable this(null, tasktype, bytesComplete, totalBytes); } - public CompactionInfo(Integer id, OperationType tasktype, long bytesComplete, long totalBytes) + public CompactionInfo(UUID id, OperationType tasktype, long bytesComplete, long totalBytes) { this.tasktype = tasktype; this.bytesComplete = bytesComplete; @@ -53,7 +54,7 @@ public final class CompactionInfo implements Serializable return new CompactionInfo(cfm == null ? null : cfm.cfId, tasktype, bytesComplete, totalBytes); } - public Integer getId() + public UUID getId() { return cfm == null ? null : cfm.cfId; } @@ -100,7 +101,7 @@ public final class CompactionInfo implements Serializable public Map asMap() { Map ret = new HashMap(); - ret.put("id", Integer.toString(getId())); + ret.put("id", getId().toString()); ret.put("keyspace", getKeyspace()); ret.put("columnfamily", getColumnFamily()); ret.put("bytesComplete", Long.toString(bytesComplete)); diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index f2cc8dbbec..68c23fd371 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -1209,8 +1209,7 @@ public class CompactionManager implements CompactionManagerMBean { try { - return new CompactionInfo(this.hashCode(), - OperationType.CLEANUP, + return new CompactionInfo(OperationType.CLEANUP, scanner.getCurrentPosition(), scanner.getLengthInBytes()); } @@ -1235,8 +1234,7 @@ public class CompactionManager implements CompactionManagerMBean { try { - return new CompactionInfo(this.hashCode(), - OperationType.SCRUB, + return new CompactionInfo(OperationType.SCRUB, dataFile.getFilePointer(), dataFile.length()); } diff --git a/src/java/org/apache/cassandra/db/index/SecondaryIndexBuilder.java b/src/java/org/apache/cassandra/db/index/SecondaryIndexBuilder.java index 69f0915c5d..fbead88fa0 100644 --- a/src/java/org/apache/cassandra/db/index/SecondaryIndexBuilder.java +++ b/src/java/org/apache/cassandra/db/index/SecondaryIndexBuilder.java @@ -47,8 +47,7 @@ public class SecondaryIndexBuilder extends CompactionInfo.Holder public CompactionInfo getCompactionInfo() { - return new CompactionInfo(this.hashCode(), - OperationType.INDEX_BUILD, + return new CompactionInfo(OperationType.INDEX_BUILD, iter.getBytesRead(), iter.getTotalBytes()); } diff --git a/src/java/org/apache/cassandra/streaming/StreamRequest.java b/src/java/org/apache/cassandra/streaming/StreamRequest.java index 2b6269284d..a8de4a663b 100644 --- a/src/java/org/apache/cassandra/streaming/StreamRequest.java +++ b/src/java/org/apache/cassandra/streaming/StreamRequest.java @@ -27,6 +27,7 @@ import java.util.List; import com.google.common.collect.Iterables; +import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.db.Table; @@ -136,7 +137,7 @@ public class StreamRequest dos.writeInt(Iterables.size(srm.columnFamilies)); for (ColumnFamilyStore cfs : srm.columnFamilies) - dos.writeInt(cfs.metadata.cfId); + ColumnFamily.serializer.serializeCfId(cfs.metadata.cfId, dos, version); } } @@ -163,7 +164,7 @@ public class StreamRequest List stores = new ArrayList(); int cfsSize = dis.readInt(); for (int i = 0; i < cfsSize; ++i) - stores.add(Table.open(table).getColumnFamilyStore(dis.readInt())); + stores.add(Table.open(table).getColumnFamilyStore(ColumnFamily.serializer.deserializeCfId(dis, version))); return new StreamRequest(target, ranges, table, stores, sessionId, type); } @@ -184,7 +185,7 @@ public class StreamRequest size += TypeSizes.NATIVE.sizeof(sr.type.name()); size += TypeSizes.NATIVE.sizeof(Iterables.size(sr.columnFamilies)); for (ColumnFamilyStore cfs : sr.columnFamilies) - size += TypeSizes.NATIVE.sizeof(cfs.metadata.cfId); + size += ColumnFamily.serializer.cfIdSerializedSize(cfs.metadata.cfId, TypeSizes.NATIVE, version); return size; } } diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 577f3ef840..15e493f4c9 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -606,7 +606,7 @@ public class FBUtilities DataOutputBuffer buffer = new DataOutputBuffer(size); serializer.serialize(object, buffer, version); assert buffer.getLength() == size && buffer.getData().length == size - : String.format("Final buffer length %s to accomodate data size of %s (predicted %s) for %s", + : String.format("Final buffer length %s to accommodate data size of %s (predicted %s) for %s", buffer.getData().length, buffer.getLength(), size, object); return buffer.getData(); } diff --git a/test/data/serialization/1.2/db.Row.bin b/test/data/serialization/1.2/db.Row.bin index 121bd457d7957746584f65eae29417631f6ced80..3b080a921d290004adcbf7679c5403c46de6462c 100644 GIT binary patch delta 84 zcmaFQ+|Q!Gz>=L>X~5`fEaUP)!a#fCA>H#+j_hAs|NsC01_mIQC@_n8|6`enYn>$D f{g)HuSYfqw;+hX?!8_BE5o&&H0;!q2meC0SueBs; delta 53 zcmeBYdCx4sz>=L>X~4+9!2F{A|Ns9D6BTAL?`rCpxYvmjDE1yKHd%sk6G(9KSw<%S DEX5XF diff --git a/test/data/serialization/1.2/db.RowMutation.bin b/test/data/serialization/1.2/db.RowMutation.bin index 5507eebfa4c343ef3fb25487c26fda480b01e205..ed0aba5ee010de19f67f81cbfaa7bfa3ccef471b 100644 GIT binary patch literal 3602 zcmZSJ^iHiTE=WvHHDq7~G7StFKmbXUB^xLT6k_x>mT~zYVW2(nknZ^@NA|D95UBtE z|9=BWHv+ITup|NjLIA{KNdf|7A#RXBl5SaIPH8FwOEM5ZWtm?98QdU&WSDFU5HLcd z!LlqsmwjPiNd*F!@JpaH7f1l6Ck+Ug7#Nse`0sx#1JTQV0i*~d0aKI?1YkuB4D6Sn ziY|f_fh3@c(1HaNR*did%L#It!WCYl1vN`OuwVui98e&JT%o~}5CfVrMsZ|i zqxBxV-a`*3ee`e|tq6gQ3*7An^zf5F4nKHu#Ha`{1crUhh#pV|*aK>`=*Ly`qXi(7 Q5!L`4tvhix@ki@U03|{_JOBUy literal 3266 zcmeHJ!4APd5S{AM&?E5|628Gv!oe4`wH%O0gv5zB`6ge1uc(uY6Ox+QcAGWf2W;PA zyM6OE&AiOc6skSBSuTP|cA$*gb_WOsE2hXI`0Zu0}&wfYx)!+!lZm6 z!qz&Ntt5yDkwx7YH$D-Oj*bxcC4!0y{Q4b)L}>bou^JKj%otH3{ouGisH5PWNJIfG znjn-2z`JCkGl?jWiCF672W5Rz3rpRB%COrH=92OMYe`-sWQS$t@R55_51&0lT`N@- z|9uy%?gV11nqk^Gac!D|BEKHMWlvLP**sxhICbttF<>O~2wtW~d0AO9ahMrS7U^$XH-$|AQ?P+e5UdMm0teKXX%pH6&~2EU%cqJ}_hcne>=FRS2soGk delta 147 zcmaEFe%*XRn+5{|^9u$Ds9~vRU|?iqFeykZ$%skGPc6<TB;#9^z*;7;uyDR`f)i7BA diff --git a/test/unit/org/apache/cassandra/SchemaLoader.java b/test/unit/org/apache/cassandra/SchemaLoader.java index 1d3bd834b6..eb4f3dcb1f 100644 --- a/test/unit/org/apache/cassandra/SchemaLoader.java +++ b/test/unit/org/apache/cassandra/SchemaLoader.java @@ -22,6 +22,7 @@ import java.io.File; import java.io.IOException; import java.nio.ByteBuffer; import java.util.*; +import java.util.concurrent.atomic.AtomicInteger; import com.google.common.base.Charsets; @@ -50,8 +51,15 @@ public class SchemaLoader { private static Logger logger = LoggerFactory.getLogger(SchemaLoader.class); + private static AtomicInteger oldCfIdGenerator = new AtomicInteger(1000); + @BeforeClass public static void loadSchema() throws IOException + { + loadSchema(false); + } + + public static void loadSchema(boolean withOldCfIds) throws IOException { // Cleanup first cleanupAndLeaveDirs(); @@ -71,7 +79,7 @@ public class SchemaLoader startGossiper(); try { - for (KSMetaData ksm : schemaDefinition()) + for (KSMetaData ksm : schemaDefinition(withOldCfIds)) MigrationManager.announceNewKeyspace(ksm); } catch (ConfigurationException e) @@ -91,7 +99,7 @@ public class SchemaLoader Gossiper.instance.stop(); } - public static Collection schemaDefinition() throws ConfigurationException + public static Collection schemaDefinition(boolean withOldCfIds) throws ConfigurationException { List schema = new ArrayList(); @@ -151,26 +159,26 @@ public class SchemaLoader opts_rf1, // Column Families - standardCFMD(ks1, "Standard1"), - standardCFMD(ks1, "Standard2"), - standardCFMD(ks1, "Standard3"), - standardCFMD(ks1, "Standard4"), - standardCFMD(ks1, "StandardLong1"), - standardCFMD(ks1, "StandardLong2"), + standardCFMD(ks1, "Standard1", withOldCfIds), + standardCFMD(ks1, "Standard2", withOldCfIds), + standardCFMD(ks1, "Standard3", withOldCfIds), + standardCFMD(ks1, "Standard4", withOldCfIds), + standardCFMD(ks1, "StandardLong1", withOldCfIds), + standardCFMD(ks1, "StandardLong2", withOldCfIds), new CFMetaData(ks1, "ValuesWithQuotes", st, BytesType.instance, null) .defaultValidator(UTF8Type.instance), - superCFMD(ks1, "Super1", LongType.instance), - superCFMD(ks1, "Super2", LongType.instance), - superCFMD(ks1, "Super3", LongType.instance), - superCFMD(ks1, "Super4", UTF8Type.instance), - superCFMD(ks1, "Super5", bytes), - superCFMD(ks1, "Super6", LexicalUUIDType.instance, UTF8Type.instance), - indexCFMD(ks1, "Indexed1", true), - indexCFMD(ks1, "Indexed2", false), + superCFMD(ks1, "Super1", LongType.instance, withOldCfIds), + superCFMD(ks1, "Super2", LongType.instance, withOldCfIds), + superCFMD(ks1, "Super3", LongType.instance, withOldCfIds), + superCFMD(ks1, "Super4", UTF8Type.instance, withOldCfIds), + superCFMD(ks1, "Super5", bytes, withOldCfIds), + superCFMD(ks1, "Super6", LexicalUUIDType.instance, UTF8Type.instance, withOldCfIds), + indexCFMD(ks1, "Indexed1", true, withOldCfIds), + indexCFMD(ks1, "Indexed2", false, withOldCfIds), new CFMetaData(ks1, "StandardInteger1", st, @@ -188,12 +196,12 @@ public class SchemaLoader bytes, bytes) .defaultValidator(CounterColumnType.instance), - superCFMD(ks1, "SuperDirectGC", BytesType.instance).gcGraceSeconds(0), - jdbcCFMD(ks1, "JdbcInteger", IntegerType.instance).columnMetadata(integerColumn), - jdbcCFMD(ks1, "JdbcUtf8", UTF8Type.instance).columnMetadata(utf8Column), - jdbcCFMD(ks1, "JdbcLong", LongType.instance), - jdbcCFMD(ks1, "JdbcBytes", bytes), - jdbcCFMD(ks1, "JdbcAscii", AsciiType.instance), + superCFMD(ks1, "SuperDirectGC", BytesType.instance, withOldCfIds).gcGraceSeconds(0), + jdbcCFMD(ks1, "JdbcInteger", IntegerType.instance, withOldCfIds).columnMetadata(integerColumn), + jdbcCFMD(ks1, "JdbcUtf8", UTF8Type.instance, withOldCfIds).columnMetadata(utf8Column), + jdbcCFMD(ks1, "JdbcLong", LongType.instance, withOldCfIds), + jdbcCFMD(ks1, "JdbcBytes", bytes, withOldCfIds), + jdbcCFMD(ks1, "JdbcAscii", AsciiType.instance, withOldCfIds), new CFMetaData(ks1, "StandardComposite", st, @@ -204,7 +212,8 @@ public class SchemaLoader st, dynamicComposite, null), - standardCFMD(ks1, "StandardLeveled").compactionStrategyClass(LeveledCompactionStrategy.class) + standardCFMD(ks1, "StandardLeveled", withOldCfIds) + .compactionStrategyClass(LeveledCompactionStrategy.class) .compactionStrategyOptions(leveledOptions))); // Keyspace 2 @@ -213,11 +222,11 @@ public class SchemaLoader opts_rf1, // Column Families - standardCFMD(ks2, "Standard1"), - standardCFMD(ks2, "Standard3"), - superCFMD(ks2, "Super3", bytes), - superCFMD(ks2, "Super4", TimeUUIDType.instance), - indexCFMD(ks2, "Indexed1", true))); + standardCFMD(ks2, "Standard1", withOldCfIds), + standardCFMD(ks2, "Standard3", withOldCfIds), + superCFMD(ks2, "Super3", bytes, withOldCfIds), + superCFMD(ks2, "Super4", TimeUUIDType.instance, withOldCfIds), + indexCFMD(ks2, "Indexed1", true, withOldCfIds))); // Keyspace 3 schema.add(KSMetaData.testMetadata(ks3, @@ -225,8 +234,8 @@ public class SchemaLoader opts_rf5, // Column Families - standardCFMD(ks3, "Standard1"), - indexCFMD(ks3, "Indexed1", true))); + standardCFMD(ks3, "Standard1", withOldCfIds), + indexCFMD(ks3, "Indexed1", true, withOldCfIds))); // Keyspace 4 schema.add(KSMetaData.testMetadata(ks4, @@ -234,10 +243,10 @@ public class SchemaLoader opts_rf3, // Column Families - standardCFMD(ks4, "Standard1"), - standardCFMD(ks4, "Standard3"), - superCFMD(ks4, "Super3", bytes), - superCFMD(ks4, "Super4", TimeUUIDType.instance), + standardCFMD(ks4, "Standard1", withOldCfIds), + standardCFMD(ks4, "Standard3", withOldCfIds), + superCFMD(ks4, "Super3", bytes, withOldCfIds), + superCFMD(ks4, "Super4", TimeUUIDType.instance, withOldCfIds), new CFMetaData(ks4, "Super5", su, @@ -248,35 +257,35 @@ public class SchemaLoader schema.add(KSMetaData.testMetadata(ks5, simple, opts_rf2, - standardCFMD(ks5, "Standard1"), - standardCFMD(ks5, "Counter1") + standardCFMD(ks5, "Standard1", withOldCfIds), + standardCFMD(ks5, "Counter1", withOldCfIds) .defaultValidator(CounterColumnType.instance))); // Keyspace 6 schema.add(KSMetaData.testMetadata(ks6, simple, opts_rf1, - indexCFMD(ks6, "Indexed1", true))); + indexCFMD(ks6, "Indexed1", true, withOldCfIds))); // KeyCacheSpace schema.add(KSMetaData.testMetadata(ks_kcs, simple, opts_rf1, - standardCFMD(ks_kcs, "Standard1"), - standardCFMD(ks_kcs, "Standard2"), - standardCFMD(ks_kcs, "Standard3"))); + standardCFMD(ks_kcs, "Standard1", withOldCfIds), + standardCFMD(ks_kcs, "Standard2", withOldCfIds), + standardCFMD(ks_kcs, "Standard3", withOldCfIds))); // RowCacheSpace schema.add(KSMetaData.testMetadata(ks_rcs, simple, opts_rf1, - standardCFMD(ks_rcs, "CFWithoutCache").caching(CFMetaData.Caching.NONE), - standardCFMD(ks_rcs, "CachedCF").caching(CFMetaData.Caching.ALL))); + standardCFMD(ks_rcs, "CFWithoutCache", withOldCfIds).caching(CFMetaData.Caching.NONE), + standardCFMD(ks_rcs, "CachedCF", withOldCfIds).caching(CFMetaData.Caching.ALL))); schema.add(KSMetaData.testMetadataNotDurable(ks_nocommit, simple, opts_rf1, - standardCFMD(ks_nocommit, "Standard1"))); + standardCFMD(ks_nocommit, "Standard1", withOldCfIds))); if (Boolean.parseBoolean(System.getProperty("cassandra.test.compression", "false"))) @@ -296,21 +305,31 @@ public class SchemaLoader } } - private static CFMetaData standardCFMD(String ksName, String cfName) + private static CFMetaData standardCFMD(String ksName, String cfName, boolean withOldCfIds) { - return new CFMetaData(ksName, cfName, ColumnFamilyType.Standard, BytesType.instance, null); + CFMetaData cfmd = new CFMetaData(ksName, cfName, ColumnFamilyType.Standard, BytesType.instance, null); + + if (withOldCfIds) + Schema.instance.addOldCfIdMapping(oldCfIdGenerator.getAndIncrement(), cfmd.cfId); + + return cfmd; } - private static CFMetaData superCFMD(String ksName, String cfName, AbstractType subcc) + private static CFMetaData superCFMD(String ksName, String cfName, AbstractType subcc, boolean withOldCfIds) { - return superCFMD(ksName, cfName, BytesType.instance, subcc); + return superCFMD(ksName, cfName, BytesType.instance, subcc, withOldCfIds); } - private static CFMetaData superCFMD(String ksName, String cfName, AbstractType cc, AbstractType subcc) + private static CFMetaData superCFMD(String ksName, String cfName, AbstractType cc, AbstractType subcc, boolean withOldCfIds) { - return new CFMetaData(ksName, cfName, ColumnFamilyType.Super, cc, subcc); + CFMetaData cfmd = new CFMetaData(ksName, cfName, ColumnFamilyType.Super, cc, subcc); + + if (withOldCfIds) + Schema.instance.addOldCfIdMapping(oldCfIdGenerator.getAndIncrement(), cfmd.cfId); + + return cfmd; } - private static CFMetaData indexCFMD(String ksName, String cfName, final Boolean withIdxType) throws ConfigurationException + private static CFMetaData indexCFMD(String ksName, String cfName, final Boolean withIdxType, boolean withOldCfIds) throws ConfigurationException { - return standardCFMD(ksName, cfName) + return standardCFMD(ksName, cfName, withOldCfIds) .keyValidator(AsciiType.instance) .columnMetadata(new HashMap() {{ @@ -319,9 +338,14 @@ public class SchemaLoader put(cName, new ColumnDefinition(cName, LongType.instance, keys, null, withIdxType ? ByteBufferUtil.bytesToHex(cName) : null, null)); }}); } - private static CFMetaData jdbcCFMD(String ksName, String cfName, AbstractType comp) + private static CFMetaData jdbcCFMD(String ksName, String cfName, AbstractType comp, boolean withOldCfIds) { - return new CFMetaData(ksName, cfName, ColumnFamilyType.Standard, comp, null).defaultValidator(comp); + CFMetaData cfmd = new CFMetaData(ksName, cfName, ColumnFamilyType.Standard, comp, null).defaultValidator(comp); + + if (withOldCfIds) + Schema.instance.addOldCfIdMapping(oldCfIdGenerator.getAndIncrement(), cfmd.cfId); + + return cfmd; } public static void cleanupAndLeaveDirs() throws IOException diff --git a/test/unit/org/apache/cassandra/Util.java b/test/unit/org/apache/cassandra/Util.java index bc67e3a02f..d53c1b2adf 100644 --- a/test/unit/org/apache/cassandra/Util.java +++ b/test/unit/org/apache/cassandra/Util.java @@ -159,7 +159,7 @@ public class Util { IMutation first = rms.get(0); String tablename = first.getTable(); - Integer cfid = first.getColumnFamilyIds().iterator().next(); + UUID cfid = first.getColumnFamilyIds().iterator().next(); for (IMutation rm : rms) rm.apply(); diff --git a/test/unit/org/apache/cassandra/cache/CacheProviderTest.java b/test/unit/org/apache/cassandra/cache/CacheProviderTest.java index 1e91da0b54..46164bff14 100644 --- a/test/unit/org/apache/cassandra/cache/CacheProviderTest.java +++ b/test/unit/org/apache/cassandra/cache/CacheProviderTest.java @@ -24,6 +24,8 @@ package org.apache.cassandra.cache; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; +import java.util.UUID; + import org.junit.Test; import org.apache.cassandra.SchemaLoader; @@ -121,15 +123,17 @@ public class CacheProviderTest extends SchemaLoader @Test public void testKeys() { + UUID cfId = UUID.randomUUID(); + byte[] b1 = {1, 2, 3, 4}; - RowCacheKey key1 = new RowCacheKey(123, ByteBuffer.wrap(b1)); + RowCacheKey key1 = new RowCacheKey(cfId, ByteBuffer.wrap(b1)); byte[] b2 = {1, 2, 3, 4}; - RowCacheKey key2 = new RowCacheKey(123, ByteBuffer.wrap(b2)); + RowCacheKey key2 = new RowCacheKey(cfId, ByteBuffer.wrap(b2)); assertEquals(key1, key2); assertEquals(key1.hashCode(), key2.hashCode()); byte[] b3 = {1, 2, 3, 5}; - RowCacheKey key3 = new RowCacheKey(123, ByteBuffer.wrap(b3)); + RowCacheKey key3 = new RowCacheKey(cfId, ByteBuffer.wrap(b3)); assertNotSame(key1, key3); assertNotSame(key1.hashCode(), key3.hashCode()); } diff --git a/test/unit/org/apache/cassandra/config/DefsTest.java b/test/unit/org/apache/cassandra/config/DefsTest.java index 2da61f9f6e..ae6ccc0508 100644 --- a/test/unit/org/apache/cassandra/config/DefsTest.java +++ b/test/unit/org/apache/cassandra/config/DefsTest.java @@ -38,8 +38,6 @@ import org.apache.cassandra.io.sstable.SSTableDeletingTask; import org.apache.cassandra.locator.OldNetworkTopologyStrategy; import org.apache.cassandra.locator.SimpleStrategy; import org.apache.cassandra.service.MigrationManager; -import org.apache.cassandra.thrift.CfDef; -import org.apache.cassandra.thrift.ColumnDef; import org.apache.cassandra.thrift.IndexType; import org.apache.cassandra.utils.ByteBufferUtil; @@ -51,11 +49,8 @@ public class DefsTest extends SchemaLoader @Test public void ensureStaticCFMIdsAreLessThan1000() { - assert CFMetaData.StatusCf.cfId == 0; - assert CFMetaData.HintsCf.cfId == 1; - assert CFMetaData.MigrationsCf.cfId == 2; - assert CFMetaData.SchemaCf.cfId == 3; - assert CFMetaData.HostIdCf.cfId == 11; + assert CFMetaData.StatusCf.cfId.equals(CFMetaData.getId(Table.SYSTEM_TABLE, SystemTable.STATUS_CF)); + assert CFMetaData.HintsCf.cfId.equals(CFMetaData.getId(Table.SYSTEM_TABLE, HintedHandOffManager.HINTS_CF)); } @Test @@ -465,7 +460,7 @@ public class DefsTest extends SchemaLoader assert Schema.instance.getCFMetaData(cf.ksName, cf.cfName).getDefaultValidator() == UTF8Type.instance; // Change cfId - newCfm = new CFMetaData(cf.ksName, cf.cfName, cf.cfType, cf.comparator, cf.subcolumnComparator, cf.cfId + 1); + newCfm = new CFMetaData(cf.ksName, cf.cfName, cf.cfType, cf.comparator, cf.subcolumnComparator, UUID.randomUUID()); CFMetaData.copyOpts(newCfm, cf); try { @@ -475,7 +470,7 @@ public class DefsTest extends SchemaLoader catch (ConfigurationException expected) {} // Change cfName - newCfm = new CFMetaData(cf.ksName, cf.cfName + "_renamed", cf.cfType, cf.comparator, cf.subcolumnComparator, cf.cfId); + newCfm = new CFMetaData(cf.ksName, cf.cfName + "_renamed", cf.cfType, cf.comparator, cf.subcolumnComparator); CFMetaData.copyOpts(newCfm, cf); try { @@ -485,7 +480,7 @@ public class DefsTest extends SchemaLoader catch (ConfigurationException expected) {} // Change ksName - newCfm = new CFMetaData(cf.ksName + "_renamed", cf.cfName, cf.cfType, cf.comparator, cf.subcolumnComparator, cf.cfId); + newCfm = new CFMetaData(cf.ksName + "_renamed", cf.cfName, cf.cfType, cf.comparator, cf.subcolumnComparator); CFMetaData.copyOpts(newCfm, cf); try { @@ -495,7 +490,7 @@ public class DefsTest extends SchemaLoader catch (ConfigurationException expected) {} // Change cf type - newCfm = new CFMetaData(cf.ksName, cf.cfName, ColumnFamilyType.Super, cf.comparator, cf.subcolumnComparator, cf.cfId); + newCfm = new CFMetaData(cf.ksName, cf.cfName, ColumnFamilyType.Super, cf.comparator, cf.subcolumnComparator); CFMetaData.copyOpts(newCfm, cf); try { @@ -505,7 +500,7 @@ public class DefsTest extends SchemaLoader catch (ConfigurationException expected) {} // Change comparator - newCfm = new CFMetaData(cf.ksName, cf.cfName, cf.cfType, TimeUUIDType.instance, cf.subcolumnComparator, cf.cfId); + newCfm = new CFMetaData(cf.ksName, cf.cfName, cf.cfType, TimeUUIDType.instance, cf.subcolumnComparator); CFMetaData.copyOpts(newCfm, cf); try { diff --git a/test/unit/org/apache/cassandra/db/CommitLogTest.java b/test/unit/org/apache/cassandra/db/CommitLogTest.java index 4e48d73a3f..5c52d07236 100644 --- a/test/unit/org/apache/cassandra/db/CommitLogTest.java +++ b/test/unit/org/apache/cassandra/db/CommitLogTest.java @@ -21,6 +21,7 @@ package org.apache.cassandra.db; import java.io.*; import java.nio.ByteBuffer; +import java.util.UUID; import java.util.zip.CRC32; import java.util.zip.Checksum; @@ -111,7 +112,7 @@ public class CommitLogTest extends SchemaLoader assert CommitLog.instance.activeSegments() == 2 : "Expecting 2 segments, got " + CommitLog.instance.activeSegments(); - int cfid2 = rm2.getColumnFamilyIds().iterator().next(); + UUID cfid2 = rm2.getColumnFamilyIds().iterator().next(); CommitLog.instance.discardCompletedSegments(cfid2, CommitLog.instance.getContext().get()); // Assert we still have both our segment @@ -133,7 +134,7 @@ public class CommitLogTest extends SchemaLoader assert CommitLog.instance.activeSegments() == 1 : "Expecting 1 segment, got " + CommitLog.instance.activeSegments(); // "Flush": this won't delete anything - int cfid1 = rm.getColumnFamilyIds().iterator().next(); + UUID cfid1 = rm.getColumnFamilyIds().iterator().next(); CommitLog.instance.discardCompletedSegments(cfid1, CommitLog.instance.getContext().get()); assert CommitLog.instance.activeSegments() == 1 : "Expecting 1 segment, got " + CommitLog.instance.activeSegments(); @@ -151,7 +152,7 @@ public class CommitLogTest extends SchemaLoader // "Flush" second cf: The first segment should be deleted since we // didn't write anything on cf1 since last flush (and we flush cf2) - int cfid2 = rm2.getColumnFamilyIds().iterator().next(); + UUID cfid2 = rm2.getColumnFamilyIds().iterator().next(); CommitLog.instance.discardCompletedSegments(cfid2, CommitLog.instance.getContext().get()); // Assert we still have both our segment diff --git a/test/unit/org/apache/cassandra/db/SerializationsTest.java b/test/unit/org/apache/cassandra/db/SerializationsTest.java index 264fd3a581..5f3e3e44f6 100644 --- a/test/unit/org/apache/cassandra/db/SerializationsTest.java +++ b/test/unit/org/apache/cassandra/db/SerializationsTest.java @@ -34,21 +34,24 @@ import org.apache.cassandra.service.StorageService; import org.apache.cassandra.thrift.SlicePredicate; import org.apache.cassandra.thrift.SliceRange; import org.apache.cassandra.utils.ByteBufferUtil; -import org.apache.cassandra.utils.FBUtilities; +import org.junit.BeforeClass; import org.junit.Test; import java.io.DataInputStream; import java.io.DataOutputStream; import java.io.IOException; import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.HashMap; +import java.util.*; public class SerializationsTest extends AbstractSerializationsTester { + @BeforeClass + public static void loadSchema() throws IOException + { + loadSchema(true); + } + private void testRangeSliceCommandWrite() throws IOException { ByteBuffer startCol = ByteBufferUtil.bytes("Start"); @@ -213,7 +216,7 @@ public class SerializationsTest extends AbstractSerializationsTester standardRm.add(Statics.StandardCf); RowMutation superRm = new RowMutation(Statics.KS, Statics.Key); superRm.add(Statics.SuperCf); - Map mods = new HashMap(); + Map mods = new HashMap(); mods.put(Statics.StandardCf.metadata().cfId, Statics.StandardCf); mods.put(Statics.SuperCf.metadata().cfId, Statics.SuperCf); RowMutation mixedRm = new RowMutation(Statics.KS, Statics.Key, mods);