mirror of https://github.com/apache/cassandra
Adding columns via ALTER TABLE can generate corrupt sstables
Patch by marcuse; reviewed by Alex Petrov and Sam Tunnicliffe for CASSANDRA-16735
This commit is contained in:
parent
9b6dd382bd
commit
1d87da3f6f
|
|
@ -1,4 +1,5 @@
|
||||||
3.0.25:
|
3.0.25:
|
||||||
|
* Adding columns via ALTER TABLE can generate corrupt sstables (CASSANDRA-16735)
|
||||||
* Add flag to disable ALTER...DROP COMPACT STORAGE statements (CASSANDRA-16733)
|
* Add flag to disable ALTER...DROP COMPACT STORAGE statements (CASSANDRA-16733)
|
||||||
* Clean transaction log leftovers at the beginning of sstablelevelreset and sstableofflinerelevel (CASSANDRA-12519)
|
* Clean transaction log leftovers at the beginning of sstablelevelreset and sstableofflinerelevel (CASSANDRA-12519)
|
||||||
* CQL shell should prefer newer TLS version by default (CASSANDRA-16695)
|
* CQL shell should prefer newer TLS version by default (CASSANDRA-16695)
|
||||||
|
|
|
||||||
|
|
@ -190,29 +190,6 @@ public class ColumnDefinition extends ColumnSpecification implements Comparable<
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
private static class Placeholder extends ColumnDefinition
|
|
||||||
{
|
|
||||||
Placeholder(CFMetaData table, ByteBuffer name, AbstractType<?> type, int position, Kind kind)
|
|
||||||
{
|
|
||||||
super(table, name, type, position, kind);
|
|
||||||
}
|
|
||||||
|
|
||||||
public boolean isPlaceholder()
|
|
||||||
{
|
|
||||||
return true;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
public static ColumnDefinition placeholder(CFMetaData table, ByteBuffer name, boolean isStatic)
|
|
||||||
{
|
|
||||||
return new Placeholder(table, name, EmptyType.instance, NO_POSITION, isStatic ? Kind.STATIC : Kind.REGULAR);
|
|
||||||
}
|
|
||||||
|
|
||||||
public boolean isPlaceholder()
|
|
||||||
{
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
|
|
||||||
public ColumnDefinition copy()
|
public ColumnDefinition copy()
|
||||||
{
|
{
|
||||||
return new ColumnDefinition(ksName, cfName, name, type, position, kind);
|
return new ColumnDefinition(ksName, cfName, name, type, position, kind);
|
||||||
|
|
|
||||||
|
|
@ -425,7 +425,7 @@ public class Columns extends AbstractCollection<ColumnDefinition> implements Col
|
||||||
return size;
|
return size;
|
||||||
}
|
}
|
||||||
|
|
||||||
public Columns deserialize(DataInputPlus in, CFMetaData metadata, boolean isStatic) throws IOException
|
public Columns deserialize(DataInputPlus in, CFMetaData metadata) throws IOException
|
||||||
{
|
{
|
||||||
int length = (int)in.readUnsignedVInt();
|
int length = (int)in.readUnsignedVInt();
|
||||||
BTree.Builder<ColumnDefinition> builder = BTree.builder(Comparator.naturalOrder());
|
BTree.Builder<ColumnDefinition> builder = BTree.builder(Comparator.naturalOrder());
|
||||||
|
|
@ -441,29 +441,14 @@ public class Columns extends AbstractCollection<ColumnDefinition> implements Col
|
||||||
// fail deserialization because of that. So we grab a "fake" ColumnDefinition that ensure proper
|
// fail deserialization because of that. So we grab a "fake" ColumnDefinition that ensure proper
|
||||||
// deserialization. The column will be ignore later on anyway.
|
// deserialization. The column will be ignore later on anyway.
|
||||||
column = metadata.getDroppedColumnDefinition(name);
|
column = metadata.getDroppedColumnDefinition(name);
|
||||||
|
|
||||||
// If there's no dropped column, it may be for a column we haven't received a schema update for yet
|
|
||||||
// so we create a placeholder column. If this is a read, the placeholder column will let the response
|
|
||||||
// serializer know we're not serializing all requested columns when it writes the row flags, but it
|
|
||||||
// will cause mutations that try to write values for this column to fail.
|
|
||||||
if (column == null)
|
if (column == null)
|
||||||
column = ColumnDefinition.placeholder(metadata, name, isStatic);
|
throw new RuntimeException("Unknown column " + UTF8Type.instance.getString(name) + " during deserialization");
|
||||||
}
|
}
|
||||||
builder.add(column);
|
builder.add(column);
|
||||||
}
|
}
|
||||||
return new Columns(builder.build());
|
return new Columns(builder.build());
|
||||||
}
|
}
|
||||||
|
|
||||||
public Columns deserializeStatics(DataInputPlus in, CFMetaData metadata) throws IOException
|
|
||||||
{
|
|
||||||
return deserialize(in, metadata, true);
|
|
||||||
}
|
|
||||||
|
|
||||||
public Columns deserializeRegulars(DataInputPlus in, CFMetaData metadata) throws IOException
|
|
||||||
{
|
|
||||||
return deserialize(in, metadata, false);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* If both ends have a pre-shared superset of the columns we are serializing, we can send them much
|
* If both ends have a pre-shared superset of the columns we are serializing, we can send them much
|
||||||
* more efficiently. Both ends must provide the identically same set of columns.
|
* more efficiently. Both ends must provide the identically same set of columns.
|
||||||
|
|
|
||||||
|
|
@ -40,7 +40,6 @@ import org.apache.cassandra.io.sstable.metadata.IMetadataComponentSerializer;
|
||||||
import org.apache.cassandra.io.util.DataInputPlus;
|
import org.apache.cassandra.io.util.DataInputPlus;
|
||||||
import org.apache.cassandra.io.util.DataOutputPlus;
|
import org.apache.cassandra.io.util.DataOutputPlus;
|
||||||
import org.apache.cassandra.utils.ByteBufferUtil;
|
import org.apache.cassandra.utils.ByteBufferUtil;
|
||||||
import org.apache.cassandra.utils.SearchIterator;
|
|
||||||
|
|
||||||
public class SerializationHeader
|
public class SerializationHeader
|
||||||
{
|
{
|
||||||
|
|
@ -161,18 +160,6 @@ public class SerializationHeader
|
||||||
return !columns.statics.isEmpty();
|
return !columns.statics.isEmpty();
|
||||||
}
|
}
|
||||||
|
|
||||||
public boolean hasAllColumns(Row row, boolean isStatic)
|
|
||||||
{
|
|
||||||
SearchIterator<ColumnDefinition, ColumnData> rowIter = row.searchIterator();
|
|
||||||
Iterable<ColumnDefinition> columns = isStatic ? columns().statics : columns().regulars;
|
|
||||||
for (ColumnDefinition column : columns)
|
|
||||||
{
|
|
||||||
if (rowIter.next(column) == null)
|
|
||||||
return false;
|
|
||||||
}
|
|
||||||
return true;
|
|
||||||
}
|
|
||||||
|
|
||||||
public boolean isForSSTable()
|
public boolean isForSSTable()
|
||||||
{
|
{
|
||||||
return isForSSTable;
|
return isForSSTable;
|
||||||
|
|
@ -455,8 +442,8 @@ public class SerializationHeader
|
||||||
Columns statics, regulars;
|
Columns statics, regulars;
|
||||||
if (selection == null)
|
if (selection == null)
|
||||||
{
|
{
|
||||||
statics = hasStatic ? Columns.serializer.deserializeStatics(in, metadata) : Columns.NONE;
|
statics = hasStatic ? Columns.serializer.deserialize(in, metadata) : Columns.NONE;
|
||||||
regulars = Columns.serializer.deserializeRegulars(in, metadata);
|
regulars = Columns.serializer.deserialize(in, metadata);
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -17,7 +17,6 @@
|
||||||
*/
|
*/
|
||||||
package org.apache.cassandra.db;
|
package org.apache.cassandra.db;
|
||||||
|
|
||||||
import java.io.IOException;
|
|
||||||
import java.nio.ByteBuffer;
|
import java.nio.ByteBuffer;
|
||||||
|
|
||||||
import org.apache.cassandra.config.CFMetaData;
|
import org.apache.cassandra.config.CFMetaData;
|
||||||
|
|
@ -28,19 +27,14 @@ import org.apache.cassandra.utils.ByteBufferUtil;
|
||||||
* Exception thrown when we read a column internally that is unknown. Note that
|
* Exception thrown when we read a column internally that is unknown. Note that
|
||||||
* this is an internal exception and is not meant to be user facing.
|
* this is an internal exception and is not meant to be user facing.
|
||||||
*/
|
*/
|
||||||
public class UnknownColumnException extends IOException
|
public class UnknownColumnException extends Exception
|
||||||
{
|
{
|
||||||
public final ByteBuffer columnName;
|
public final ByteBuffer columnName;
|
||||||
|
|
||||||
public UnknownColumnException(String ksName, String cfName, ByteBuffer columnName)
|
|
||||||
{
|
|
||||||
super(String.format("Unknown column %s in table %s.%s", stringify(columnName), ksName, cfName));
|
|
||||||
this.columnName = columnName;
|
|
||||||
}
|
|
||||||
|
|
||||||
public UnknownColumnException(CFMetaData metadata, ByteBuffer columnName)
|
public UnknownColumnException(CFMetaData metadata, ByteBuffer columnName)
|
||||||
{
|
{
|
||||||
this(metadata.ksName, metadata.cfName, columnName);
|
super(String.format("Unknown column %s in table %s.%s", stringify(columnName), metadata.ksName, metadata.cfName));
|
||||||
|
this.columnName = columnName;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static String stringify(ByteBuffer name)
|
private static String stringify(ByteBuffer name)
|
||||||
|
|
|
||||||
|
|
@ -396,8 +396,8 @@ public class ColumnFilter
|
||||||
{
|
{
|
||||||
if (version >= MessagingService.VERSION_3014)
|
if (version >= MessagingService.VERSION_3014)
|
||||||
{
|
{
|
||||||
Columns statics = Columns.serializer.deserializeStatics(in, metadata);
|
Columns statics = Columns.serializer.deserialize(in, metadata);
|
||||||
Columns regulars = Columns.serializer.deserializeRegulars(in, metadata);
|
Columns regulars = Columns.serializer.deserialize(in, metadata);
|
||||||
fetched = new PartitionColumns(statics, regulars);
|
fetched = new PartitionColumns(statics, regulars);
|
||||||
}
|
}
|
||||||
else
|
else
|
||||||
|
|
@ -408,8 +408,8 @@ public class ColumnFilter
|
||||||
|
|
||||||
if (hasSelection)
|
if (hasSelection)
|
||||||
{
|
{
|
||||||
Columns statics = Columns.serializer.deserializeStatics(in, metadata);
|
Columns statics = Columns.serializer.deserialize(in, metadata);
|
||||||
Columns regulars = Columns.serializer.deserializeRegulars(in, metadata);
|
Columns regulars = Columns.serializer.deserialize(in, metadata);
|
||||||
selection = new PartitionColumns(statics, regulars);
|
selection = new PartitionColumns(statics, regulars);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -17,7 +17,6 @@
|
||||||
*/
|
*/
|
||||||
package org.apache.cassandra.db.partitions;
|
package org.apache.cassandra.db.partitions;
|
||||||
|
|
||||||
import java.io.IOError;
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
import java.nio.ByteBuffer;
|
import java.nio.ByteBuffer;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
|
|
@ -754,12 +753,6 @@ public class PartitionUpdate extends AbstractBTreePartition
|
||||||
deletionBuilder.add((RangeTombstoneMarker)unfiltered);
|
deletionBuilder.add((RangeTombstoneMarker)unfiltered);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
catch (IOError e)
|
|
||||||
{
|
|
||||||
if (e.getCause() != null && e.getCause() instanceof UnknownColumnException)
|
|
||||||
throw (UnknownColumnException) e.getCause();
|
|
||||||
throw e;
|
|
||||||
}
|
|
||||||
|
|
||||||
MutableDeletionInfo deletionInfo = deletionBuilder.build();
|
MutableDeletionInfo deletionInfo = deletionBuilder.build();
|
||||||
return new PartitionUpdate(metadata,
|
return new PartitionUpdate(metadata,
|
||||||
|
|
|
||||||
|
|
@ -19,9 +19,9 @@ package org.apache.cassandra.db.rows;
|
||||||
|
|
||||||
import java.io.IOException;
|
import java.io.IOException;
|
||||||
|
|
||||||
|
import com.google.common.collect.Collections2;
|
||||||
|
|
||||||
import org.apache.cassandra.config.ColumnDefinition;
|
import org.apache.cassandra.config.ColumnDefinition;
|
||||||
import org.apache.cassandra.db.marshal.UTF8Type;
|
|
||||||
import org.apache.cassandra.db.*;
|
import org.apache.cassandra.db.*;
|
||||||
import org.apache.cassandra.io.util.DataInputPlus;
|
import org.apache.cassandra.io.util.DataInputPlus;
|
||||||
import org.apache.cassandra.io.util.DataOutputPlus;
|
import org.apache.cassandra.io.util.DataOutputPlus;
|
||||||
|
|
@ -133,7 +133,7 @@ public class UnfilteredSerializer
|
||||||
LivenessInfo pkLiveness = row.primaryKeyLivenessInfo();
|
LivenessInfo pkLiveness = row.primaryKeyLivenessInfo();
|
||||||
Row.Deletion deletion = row.deletion();
|
Row.Deletion deletion = row.deletion();
|
||||||
boolean hasComplexDeletion = row.hasComplexDeletion();
|
boolean hasComplexDeletion = row.hasComplexDeletion();
|
||||||
boolean hasAllColumns = header.hasAllColumns(row, isStatic);
|
boolean hasAllColumns = (row.columnCount() == headerColumns.size());
|
||||||
boolean hasExtendedFlags = hasExtendedFlags(row);
|
boolean hasExtendedFlags = hasExtendedFlags(row);
|
||||||
|
|
||||||
if (isStatic)
|
if (isStatic)
|
||||||
|
|
@ -192,12 +192,7 @@ public class UnfilteredSerializer
|
||||||
// with. So we use the ColumnDefinition from the "header" which is "current". Also see #11810 for what
|
// with. So we use the ColumnDefinition from the "header" which is "current". Also see #11810 for what
|
||||||
// happens if we don't do that.
|
// happens if we don't do that.
|
||||||
ColumnDefinition column = si.next(data.column());
|
ColumnDefinition column = si.next(data.column());
|
||||||
|
assert column != null;
|
||||||
// we may have columns that the remote node isn't aware of due to inflight schema changes
|
|
||||||
// in cases where it tries to fetch all columns, it will set the `all columns` flag, but only
|
|
||||||
// expect a subset of columns (from this node's perspective). See CASSANDRA-15899
|
|
||||||
if (column == null)
|
|
||||||
continue;
|
|
||||||
|
|
||||||
if (data.column.isSimple())
|
if (data.column.isSimple())
|
||||||
Cell.serializer.serialize((Cell) data, column, out, pkLiveness, header);
|
Cell.serializer.serialize((Cell) data, column, out, pkLiveness, header);
|
||||||
|
|
@ -279,7 +274,7 @@ public class UnfilteredSerializer
|
||||||
LivenessInfo pkLiveness = row.primaryKeyLivenessInfo();
|
LivenessInfo pkLiveness = row.primaryKeyLivenessInfo();
|
||||||
Row.Deletion deletion = row.deletion();
|
Row.Deletion deletion = row.deletion();
|
||||||
boolean hasComplexDeletion = row.hasComplexDeletion();
|
boolean hasComplexDeletion = row.hasComplexDeletion();
|
||||||
boolean hasAllColumns = header.hasAllColumns(row, isStatic);
|
boolean hasAllColumns = (row.columnCount() == headerColumns.size());
|
||||||
|
|
||||||
if (!pkLiveness.isEmpty())
|
if (!pkLiveness.isEmpty())
|
||||||
size += header.timestampSerializedSize(pkLiveness.timestamp());
|
size += header.timestampSerializedSize(pkLiveness.timestamp());
|
||||||
|
|
@ -298,8 +293,7 @@ public class UnfilteredSerializer
|
||||||
for (ColumnData data : row)
|
for (ColumnData data : row)
|
||||||
{
|
{
|
||||||
ColumnDefinition column = si.next(data.column());
|
ColumnDefinition column = si.next(data.column());
|
||||||
if (column == null)
|
assert column != null;
|
||||||
continue;
|
|
||||||
|
|
||||||
if (data.column.isSimple())
|
if (data.column.isSimple())
|
||||||
size += Cell.serializer.serializedSize((Cell) data, column, pkLiveness, header);
|
size += Cell.serializer.serializedSize((Cell) data, column, pkLiveness, header);
|
||||||
|
|
@ -490,9 +484,6 @@ public class UnfilteredSerializer
|
||||||
Columns columns = hasAllColumns ? headerColumns : Columns.serializer.deserializeSubset(headerColumns, in);
|
Columns columns = hasAllColumns ? headerColumns : Columns.serializer.deserializeSubset(headerColumns, in);
|
||||||
for (ColumnDefinition column : columns)
|
for (ColumnDefinition column : columns)
|
||||||
{
|
{
|
||||||
// if the column is a placeholder, then it's not part of our schema, and we can't deserialize it
|
|
||||||
if (column.isPlaceholder())
|
|
||||||
throw new UnknownColumnException(column.ksName, column.cfName, column.name.bytes);
|
|
||||||
if (column.isSimple())
|
if (column.isSimple())
|
||||||
readSimpleColumn(column, in, header, helper, builder, rowLiveness);
|
readSimpleColumn(column, in, header, helper, builder, rowLiveness);
|
||||||
else
|
else
|
||||||
|
|
|
||||||
|
|
@ -18,100 +18,72 @@
|
||||||
|
|
||||||
package org.apache.cassandra.distributed.test;
|
package org.apache.cassandra.distributed.test;
|
||||||
|
|
||||||
import java.util.function.Consumer;
|
|
||||||
|
|
||||||
import org.junit.Test;
|
import org.junit.Test;
|
||||||
|
|
||||||
import org.apache.cassandra.dht.ByteOrderedPartitioner;
|
|
||||||
import org.apache.cassandra.distributed.Cluster;
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
||||||
import org.apache.cassandra.distributed.api.IInstanceConfig;
|
|
||||||
|
|
||||||
import static org.apache.cassandra.distributed.shared.AssertUtils.assertRows;
|
import static org.junit.Assert.assertTrue;
|
||||||
|
|
||||||
public class SchemaTest extends TestBaseImpl
|
public class SchemaTest extends TestBaseImpl
|
||||||
{
|
{
|
||||||
private static final Consumer<IInstanceConfig> CONFIG_CONSUMER = config -> {
|
|
||||||
config.set("partitioner", ByteOrderedPartitioner.class.getSimpleName());
|
|
||||||
config.set("initial_token", Integer.toString(config.num() * 1000));
|
|
||||||
};
|
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void dropColumnMixedMode() throws Throwable
|
public void readRepair() throws Throwable
|
||||||
{
|
{
|
||||||
try (Cluster cluster = init(Cluster.create(2, CONFIG_CONSUMER)))
|
try (Cluster cluster = init(Cluster.build(2).start()))
|
||||||
{
|
{
|
||||||
cluster.schemaChange("CREATE TABLE "+KEYSPACE+".tbl (id int primary key, v1 int, v2 int, v3 int)");
|
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v1 int, v2 int, primary key (pk, ck))");
|
||||||
Object [][] someExpected = new Object[5][];
|
String name = "aaa";
|
||||||
Object [][] allExpected1 = new Object[5][];
|
cluster.get(1).schemaChangeInternal("ALTER TABLE " + KEYSPACE + ".tbl ADD " + name + " list<int>");
|
||||||
Object [][] allExpected2 = new Object[5][];
|
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) values (?,1,1,1)", 1);
|
||||||
for (int i = 0; i < 5; i++)
|
selectSilent(cluster, name);
|
||||||
{
|
|
||||||
int v1 = i * 10, v2 = i * 100, v3 = i * 1000;
|
cluster.get(2).flush(KEYSPACE);
|
||||||
cluster.coordinator(1).execute("INSERT INTO "+KEYSPACE+".tbl (id, v1, v2, v3) VALUES (?,?,?, ?)" , ConsistencyLevel.ALL, i, v1, v2, v3);
|
cluster.get(2).schemaChangeInternal("ALTER TABLE " + KEYSPACE + ".tbl ADD " + name + " list<int>");
|
||||||
someExpected[i] = new Object[] {i, v1};
|
cluster.get(2).shutdown();
|
||||||
allExpected1[i] = new Object[] {i, v1, v3};
|
cluster.get(2).startup();
|
||||||
allExpected2[i] = new Object[] {i, v1, v2, v3};
|
cluster.get(2).forceCompact(KEYSPACE, "tbl");
|
||||||
}
|
|
||||||
cluster.forEach((instance) -> instance.flush(KEYSPACE));
|
|
||||||
cluster.get(1).schemaChangeInternal("ALTER TABLE "+KEYSPACE+".tbl DROP v2");
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT id, v1 FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), someExpected);
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), allExpected1);
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT id, v1 FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), someExpected);
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), allExpected2);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void addColumnMixedMode() throws Throwable
|
public void readRepairWithCompaction() throws Throwable
|
||||||
{
|
{
|
||||||
try (Cluster cluster = init(Cluster.create(2, CONFIG_CONSUMER)))
|
try (Cluster cluster = init(Cluster.build(2).start()))
|
||||||
{
|
{
|
||||||
cluster.schemaChange("CREATE TABLE "+KEYSPACE+".tbl (id int primary key, v1 int, v2 int)");
|
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v1 int, v2 int, primary key (pk, ck))");
|
||||||
Object [][] someExpected = new Object[5][];
|
String name = "v10";
|
||||||
Object [][] allExpected1 = new Object[5][];
|
cluster.get(1).schemaChangeInternal("ALTER TABLE " + KEYSPACE + ".tbl ADD " + name + " list<int>");
|
||||||
Object [][] allExpected2 = new Object[5][];
|
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) values (?,1,1,1)", 1);
|
||||||
for (int i = 0; i < 5; i++)
|
selectSilent(cluster, name);
|
||||||
{
|
cluster.get(2).flush(KEYSPACE);
|
||||||
int v1 = i * 10, v2 = i * 100;
|
cluster.get(2).schemaChangeInternal("ALTER TABLE " + KEYSPACE + ".tbl ADD " + name + " list<int>");
|
||||||
cluster.coordinator(1).execute("INSERT INTO "+KEYSPACE+".tbl (id, v1, v2) VALUES (?,?,?)" , ConsistencyLevel.ALL, i, v1, v2);
|
cluster.get(2).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2, " + name + ") values (?,1,1,1,[1])", 1);
|
||||||
someExpected[i] = new Object[] {i, v1};
|
cluster.get(2).flush(KEYSPACE);
|
||||||
allExpected1[i] = new Object[] {i, v1, v2, null};
|
cluster.get(2).forceCompact(KEYSPACE, "tbl");
|
||||||
allExpected2[i] = new Object[] {i, v1, v2};
|
cluster.get(2).shutdown();
|
||||||
}
|
cluster.get(2).startup();
|
||||||
cluster.forEach((instance) -> instance.flush(KEYSPACE));
|
cluster.get(2).forceCompact(KEYSPACE, "tbl");
|
||||||
cluster.get(1).schemaChangeInternal("ALTER TABLE "+KEYSPACE+".tbl ADD v3 int");
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT id, v1 FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), someExpected);
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), allExpected1);
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT id, v1 FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), someExpected);
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), allExpected2);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
private void selectSilent(Cluster cluster, String name)
|
||||||
public void addDropColumnMixedMode() throws Throwable
|
|
||||||
{
|
{
|
||||||
try (Cluster cluster = init(Cluster.create(2, CONFIG_CONSUMER)))
|
try
|
||||||
{
|
{
|
||||||
cluster.schemaChange("CREATE TABLE "+KEYSPACE+".tbl (id int primary key, v1 int, v2 int)");
|
cluster.coordinator(1).execute(withKeyspace("SELECT * FROM %s.tbl WHERE pk = ?"), ConsistencyLevel.ALL, 1);
|
||||||
Object [][] someExpected = new Object[5][];
|
}
|
||||||
Object [][] allExpected1 = new Object[5][];
|
catch (Exception e)
|
||||||
Object [][] allExpected2 = new Object[5][];
|
{
|
||||||
for (int i = 0; i < 5; i++)
|
boolean causeIsUnknownColumn = false;
|
||||||
|
Throwable cause = e;
|
||||||
|
while (cause != null)
|
||||||
{
|
{
|
||||||
int v1 = i * 10, v2 = i * 100;
|
if (cause.getMessage() != null && cause.getMessage().contains("Unknown column "+name+" during deserialization"))
|
||||||
cluster.coordinator(1).execute("INSERT INTO "+KEYSPACE+".tbl (id, v1, v2) VALUES (?,?,?)" , ConsistencyLevel.ALL, i, v1, v2);
|
causeIsUnknownColumn = true;
|
||||||
someExpected[i] = new Object[] {i, v1};
|
cause = cause.getCause();
|
||||||
allExpected1[i] = new Object[] {i, v1, v2, null};
|
|
||||||
allExpected2[i] = new Object[] {i, v1};
|
|
||||||
}
|
}
|
||||||
cluster.forEach((instance) -> instance.flush(KEYSPACE));
|
assertTrue(causeIsUnknownColumn);
|
||||||
cluster.get(1).schemaChangeInternal("ALTER TABLE "+KEYSPACE+".tbl ADD v3 int");
|
|
||||||
cluster.get(2).schemaChangeInternal("ALTER TABLE "+KEYSPACE+".tbl DROP v2");
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT id, v1 FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), someExpected);
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), allExpected1);
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT id, v1 FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), someExpected);
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT * FROM "+KEYSPACE+".tbl", ConsistencyLevel.ALL), allExpected2);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -110,9 +110,6 @@ public class SimpleReadWriteTest extends SharedClusterTestBase
|
||||||
row(1, 1, 1));
|
row(1, 1, 1));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* If a node receives a mutation for a column it's not aware of, it should fail, since it can't write the data.
|
|
||||||
*/
|
|
||||||
@Test
|
@Test
|
||||||
public void writeWithSchemaDisagreement() throws Throwable
|
public void writeWithSchemaDisagreement() throws Throwable
|
||||||
{
|
{
|
||||||
|
|
@ -129,7 +126,7 @@ public class SimpleReadWriteTest extends SharedClusterTestBase
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
cluster.coordinator(1).execute("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) VALUES (2, 2, 2, 2)",
|
cluster.coordinator(1).execute("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) VALUES (2, 2, 2, 2)",
|
||||||
ConsistencyLevel.ALL);
|
ConsistencyLevel.QUORUM);
|
||||||
}
|
}
|
||||||
catch (RuntimeException e)
|
catch (RuntimeException e)
|
||||||
{
|
{
|
||||||
|
|
@ -137,64 +134,9 @@ public class SimpleReadWriteTest extends SharedClusterTestBase
|
||||||
}
|
}
|
||||||
|
|
||||||
Assert.assertTrue(thrown.getMessage().contains("Exception occurred on node"));
|
Assert.assertTrue(thrown.getMessage().contains("Exception occurred on node"));
|
||||||
|
Assert.assertTrue(thrown.getCause().getCause().getCause().getMessage().contains("Unknown column v2 during deserialization"));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* If a node receives a mutation for a column it knows has been dropped, the write should succeed
|
|
||||||
*/
|
|
||||||
@Test
|
|
||||||
public void writeWithSchemaDisagreement2() throws Throwable
|
|
||||||
{
|
|
||||||
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v1 int, v2 int, PRIMARY KEY (pk, ck))");
|
|
||||||
|
|
||||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) VALUES (1, 1, 1, 1)");
|
|
||||||
cluster.get(2).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) VALUES (1, 1, 1, 1)");
|
|
||||||
cluster.get(3).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) VALUES (1, 1, 1, 1)");
|
|
||||||
|
|
||||||
for (int i=0; i<cluster.size(); i++)
|
|
||||||
cluster.get(i+1).flush(KEYSPACE);;
|
|
||||||
|
|
||||||
// Introduce schema disagreement
|
|
||||||
cluster.schemaChange("ALTER TABLE " + KEYSPACE + ".tbl DROP v2", 1);
|
|
||||||
|
|
||||||
// execute a write including the dropped column where the coordinator is not yet aware of the drop
|
|
||||||
// all nodes should process this without error
|
|
||||||
cluster.coordinator(2).execute("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1, v2) VALUES (2, 2, 2, 2)",
|
|
||||||
ConsistencyLevel.ALL);
|
|
||||||
// and flushing should also be fine
|
|
||||||
for (int i=0; i<cluster.size(); i++)
|
|
||||||
cluster.get(i+1).flush(KEYSPACE);;
|
|
||||||
// the results of reads will vary depending on whether the coordinator has seen the schema change
|
|
||||||
// note: read repairs will propagate the v2 value to node1, but this is safe and handled correctly
|
|
||||||
assertRows(cluster.coordinator(2).execute("SELECT * FROM " + KEYSPACE + ".tbl", ConsistencyLevel.ALL),
|
|
||||||
rows(row(1,1,1,1), row(2,2,2,2)));
|
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT * FROM " + KEYSPACE + ".tbl", ConsistencyLevel.ALL),
|
|
||||||
rows(row(1,1,1), row(2,2,2)));
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* If a node isn't aware of a column, but receives a mutation without that column, the write should succeed
|
|
||||||
*/
|
|
||||||
@Test
|
|
||||||
public void writeWithInconsequentialSchemaDisagreement() throws Throwable
|
|
||||||
{
|
|
||||||
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v1 int, PRIMARY KEY (pk, ck))");
|
|
||||||
|
|
||||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1) VALUES (1, 1, 1)");
|
|
||||||
cluster.get(2).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1) VALUES (1, 1, 1)");
|
|
||||||
cluster.get(3).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1) VALUES (1, 1, 1)");
|
|
||||||
|
|
||||||
// Introduce schema disagreement
|
|
||||||
cluster.schemaChange("ALTER TABLE " + KEYSPACE + ".tbl ADD v2 int", 1);
|
|
||||||
|
|
||||||
// this write shouldn't cause any problems because it doesn't write to the new column
|
|
||||||
cluster.coordinator(1).execute("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1) VALUES (2, 2, 2)",
|
|
||||||
ConsistencyLevel.ALL);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* If a node receives a read for a column it's not aware of, it shouldn't complain, since it won't have any data for that column
|
|
||||||
*/
|
|
||||||
@Test
|
@Test
|
||||||
public void readWithSchemaDisagreement() throws Throwable
|
public void readWithSchemaDisagreement() throws Throwable
|
||||||
{
|
{
|
||||||
|
|
@ -207,8 +149,20 @@ public class SimpleReadWriteTest extends SharedClusterTestBase
|
||||||
// Introduce schema disagreement
|
// Introduce schema disagreement
|
||||||
cluster.schemaChange("ALTER TABLE " + KEYSPACE + ".tbl ADD v2 int", 1);
|
cluster.schemaChange("ALTER TABLE " + KEYSPACE + ".tbl ADD v2 int", 1);
|
||||||
|
|
||||||
Object[][] expected = new Object[][]{new Object[]{1, 1, 1, null}};
|
Exception thrown = null;
|
||||||
assertRows(cluster.coordinator(1).execute("SELECT * FROM " + KEYSPACE + ".tbl WHERE pk = 1", ConsistencyLevel.ALL), expected);
|
try
|
||||||
|
{
|
||||||
|
assertRows(cluster.coordinator(1).execute("SELECT * FROM " + KEYSPACE + ".tbl WHERE pk = 1",
|
||||||
|
ConsistencyLevel.ALL),
|
||||||
|
row(1, 1, 1, null));
|
||||||
|
}
|
||||||
|
catch (Exception e)
|
||||||
|
{
|
||||||
|
thrown = e;
|
||||||
|
}
|
||||||
|
|
||||||
|
Assert.assertTrue(thrown.getMessage().contains("Exception occurred on node"));
|
||||||
|
Assert.assertTrue(thrown.getCause().getCause().getCause().getMessage().contains("Unknown column v2 during deserialization"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|
@ -438,12 +392,4 @@ public class SimpleReadWriteTest extends SharedClusterTestBase
|
||||||
{
|
{
|
||||||
return instance.callOnInstance(() -> Keyspace.open(KEYSPACE).getColumnFamilyStore("tbl").metric.readLatency.latency.getCount());
|
return instance.callOnInstance(() -> Keyspace.open(KEYSPACE).getColumnFamilyStore("tbl").metric.readLatency.latency.getCount());
|
||||||
}
|
}
|
||||||
|
|
||||||
private static Object[][] rows(Object[]...rows)
|
|
||||||
{
|
|
||||||
Object[][] r = new Object[rows.length][];
|
|
||||||
System.arraycopy(rows, 0, r, 0, rows.length);
|
|
||||||
return r;
|
|
||||||
}
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -132,7 +132,7 @@ public class ColumnsTest
|
||||||
{
|
{
|
||||||
Columns.serializer.serialize(columns, out);
|
Columns.serializer.serialize(columns, out);
|
||||||
Assert.assertEquals(Columns.serializer.serializedSize(columns), out.buffer().remaining());
|
Assert.assertEquals(Columns.serializer.serializedSize(columns), out.buffer().remaining());
|
||||||
Columns deserialized = Columns.serializer.deserializeRegulars(new DataInputBuffer(out.buffer(), false), mock(columns));
|
Columns deserialized = Columns.serializer.deserialize(new DataInputBuffer(out.buffer(), false), mock(columns));
|
||||||
Assert.assertEquals(columns, deserialized);
|
Assert.assertEquals(columns, deserialized);
|
||||||
Assert.assertEquals(columns.hashCode(), deserialized.hashCode());
|
Assert.assertEquals(columns.hashCode(), deserialized.hashCode());
|
||||||
assertContents(deserialized, definitions);
|
assertContents(deserialized, definitions);
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue