mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-1.2' into cassandra-2.0
Conflicts: src/java/org/apache/cassandra/streaming/StreamOut.java test/unit/org/apache/cassandra/streaming/StreamingTransferTest.java
This commit is contained in:
commit
9495eb59c4
|
|
@ -13,6 +13,7 @@ Merged from 1.2:
|
|||
* Add snitch, schema version, cluster, partitioner to JMX (CASSANDRA-5881)
|
||||
* Allow disabling SlabAllocator (CASSANDRA-5935)
|
||||
* Make user-defined compaction JMX blocking (CASSANDRA-4952)
|
||||
* Fix streaming does not transfer wrapped range (CASSANDRA-5948)
|
||||
|
||||
|
||||
2.0.0
|
||||
|
|
|
|||
|
|
@ -243,16 +243,17 @@ public class StreamSession implements IEndpointStateChangeSubscriber, IFailureDe
|
|||
if (flushTables)
|
||||
flushSSTables(stores);
|
||||
|
||||
List<Range<Token>> normalizedRanges = Range.normalize(ranges);
|
||||
List<SSTableReader> sstables = Lists.newLinkedList();
|
||||
for (ColumnFamilyStore cfStore : stores)
|
||||
{
|
||||
List<AbstractBounds<RowPosition>> rowBoundsList = Lists.newLinkedList();
|
||||
for (Range<Token> range : ranges)
|
||||
for (Range<Token> range : normalizedRanges)
|
||||
rowBoundsList.add(range.toRowBounds());
|
||||
ColumnFamilyStore.ViewFragment view = cfStore.markReferenced(rowBoundsList);
|
||||
sstables.addAll(view.sstables);
|
||||
}
|
||||
addTransferFiles(ranges, sstables);
|
||||
addTransferFiles(normalizedRanges, sstables);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -128,7 +128,7 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
* Create and transfer a single sstable, and return the keys that should have been transferred.
|
||||
* The Mutator must create the given column, but it may also create any other columns it pleases.
|
||||
*/
|
||||
private List<String> createAndTransfer(ColumnFamilyStore cfs, Mutator mutator) throws Exception
|
||||
private List<String> createAndTransfer(ColumnFamilyStore cfs, Mutator mutator, boolean transferSSTables) throws Exception
|
||||
{
|
||||
// write a temporary SSTable, and unregister it
|
||||
logger.debug("Mutating " + cfs.name);
|
||||
|
|
@ -138,18 +138,29 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
cfs.forceBlockingFlush();
|
||||
Util.compactAll(cfs).get();
|
||||
assertEquals(1, cfs.getSSTables().size());
|
||||
SSTableReader sstable = cfs.getSSTables().iterator().next();
|
||||
cfs.clearUnsafe();
|
||||
|
||||
// transfer the first and last key
|
||||
logger.debug("Transferring " + cfs.name);
|
||||
transfer(sstable);
|
||||
int[] offs;
|
||||
if (transferSSTables)
|
||||
{
|
||||
SSTableReader sstable = cfs.getSSTables().iterator().next();
|
||||
cfs.clearUnsafe();
|
||||
transferSSTables(sstable);
|
||||
offs = new int[]{1, 3};
|
||||
}
|
||||
else
|
||||
{
|
||||
long beforeStreaming = System.currentTimeMillis();
|
||||
transferRanges(cfs);
|
||||
cfs.discardSSTables(beforeStreaming);
|
||||
offs = new int[]{2, 3};
|
||||
}
|
||||
|
||||
// confirm that a single SSTable was transferred and registered
|
||||
assertEquals(1, cfs.getSSTables().size());
|
||||
|
||||
// and that the index and filter were properly recovered
|
||||
int[] offs = new int[]{1, 3};
|
||||
List<Row> rows = Util.getRangeSlice(cfs);
|
||||
assertEquals(offs.length, rows.size());
|
||||
for (int i = 0; i < offs.length; i++)
|
||||
|
|
@ -172,7 +183,7 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
return keys;
|
||||
}
|
||||
|
||||
private void transfer(SSTableReader sstable) throws Exception
|
||||
private void transferSSTables(SSTableReader sstable) throws Exception
|
||||
{
|
||||
IPartitioner p = StorageService.getPartitioner();
|
||||
List<Range<Token>> ranges = new ArrayList<>();
|
||||
|
|
@ -181,6 +192,15 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
transfer(sstable, ranges);
|
||||
}
|
||||
|
||||
private void transferRanges(ColumnFamilyStore cfs) throws Exception
|
||||
{
|
||||
IPartitioner p = StorageService.getPartitioner();
|
||||
List<Range<Token>> ranges = new ArrayList<>();
|
||||
// wrapped range
|
||||
ranges.add(new Range<Token>(p.getToken(ByteBufferUtil.bytes("key1")), p.getToken(ByteBufferUtil.bytes("key0"))));
|
||||
new StreamPlan("StreamingTransferTest").transferRanges(LOCAL, cfs.keyspace.getName(), ranges, cfs.getColumnFamilyName()).execute().get();
|
||||
}
|
||||
|
||||
private void transfer(SSTableReader sstable, List<Range<Token>> ranges) throws Exception
|
||||
{
|
||||
new StreamPlan("StreamingTransferTest").transferFiles(LOCAL, makeStreamingDetails(ranges, Arrays.asList(sstable))).execute().get();
|
||||
|
|
@ -198,6 +218,41 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
return details;
|
||||
}
|
||||
|
||||
private void doTransferTable(boolean transferSSTables) throws Exception
|
||||
{
|
||||
final Keyspace keyspace = Keyspace.open("Keyspace1");
|
||||
final ColumnFamilyStore cfs = keyspace.getColumnFamilyStore("Indexed1");
|
||||
|
||||
List<String> keys = createAndTransfer(cfs, new Mutator()
|
||||
{
|
||||
public void mutate(String key, String col, long timestamp) throws Exception
|
||||
{
|
||||
long val = key.hashCode();
|
||||
ColumnFamily cf = TreeMapBackedSortedColumns.factory.create(keyspace.getName(), cfs.name);
|
||||
cf.addColumn(column(col, "v", timestamp));
|
||||
cf.addColumn(new Column(ByteBufferUtil.bytes("birthdate"), ByteBufferUtil.bytes(val), timestamp));
|
||||
RowMutation rm = new RowMutation("Keyspace1", ByteBufferUtil.bytes(key), cf);
|
||||
logger.debug("Applying row to transfer " + rm);
|
||||
rm.apply();
|
||||
}
|
||||
}, transferSSTables);
|
||||
|
||||
// confirm that the secondary index was recovered
|
||||
for (String key : keys)
|
||||
{
|
||||
long val = key.hashCode();
|
||||
IndexExpression expr = new IndexExpression(ByteBufferUtil.bytes("birthdate"),
|
||||
IndexOperator.EQ,
|
||||
ByteBufferUtil.bytes(val));
|
||||
List<IndexExpression> clause = Arrays.asList(expr);
|
||||
IDiskAtomFilter filter = new IdentityQueryFilter();
|
||||
Range<RowPosition> range = Util.range("", "");
|
||||
List<Row> rows = cfs.search(range, clause, filter, 100);
|
||||
assertEquals(1, rows.size());
|
||||
assert rows.get(0).key.key.equals(ByteBufferUtil.bytes(key));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Test to make sure RangeTombstones at column index boundary transferred correctly.
|
||||
*/
|
||||
|
|
@ -223,7 +278,7 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
|
||||
SSTableReader sstable = cfs.getSSTables().iterator().next();
|
||||
cfs.clearUnsafe();
|
||||
transfer(sstable);
|
||||
transferSSTables(sstable);
|
||||
|
||||
// confirm that a single SSTable was transferred and registered
|
||||
assertEquals(1, cfs.getSSTables().size());
|
||||
|
|
@ -233,39 +288,15 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
}
|
||||
|
||||
@Test
|
||||
public void testTransferTable() throws Exception
|
||||
public void testTransferTableViaRanges() throws Exception
|
||||
{
|
||||
final Keyspace keyspace = Keyspace.open("Keyspace1");
|
||||
final ColumnFamilyStore cfs = keyspace.getColumnFamilyStore("Indexed1");
|
||||
doTransferTable(false);
|
||||
}
|
||||
|
||||
List<String> keys = createAndTransfer(cfs, new Mutator()
|
||||
{
|
||||
public void mutate(String key, String col, long timestamp) throws Exception
|
||||
{
|
||||
long val = key.hashCode();
|
||||
ColumnFamily cf = TreeMapBackedSortedColumns.factory.create(keyspace.getName(), cfs.name);
|
||||
cf.addColumn(column(col, "v", timestamp));
|
||||
cf.addColumn(new Column(ByteBufferUtil.bytes("birthdate"), ByteBufferUtil.bytes(val), timestamp));
|
||||
RowMutation rm = new RowMutation("Keyspace1", ByteBufferUtil.bytes(key), cf);
|
||||
logger.debug("Applying row to transfer " + rm);
|
||||
rm.apply();
|
||||
}
|
||||
});
|
||||
|
||||
// confirm that the secondary index was recovered
|
||||
for (String key : keys)
|
||||
{
|
||||
long val = key.hashCode();
|
||||
IndexExpression expr = new IndexExpression(ByteBufferUtil.bytes("birthdate"),
|
||||
IndexOperator.EQ,
|
||||
ByteBufferUtil.bytes(val));
|
||||
List<IndexExpression> clause = Arrays.asList(expr);
|
||||
IDiskAtomFilter filter = new IdentityQueryFilter();
|
||||
Range<RowPosition> range = Util.range("", "");
|
||||
List<Row> rows = cfs.search(range, clause, filter, 100);
|
||||
assertEquals(1, rows.size());
|
||||
assert rows.get(0).key.key.equals(ByteBufferUtil.bytes(key));
|
||||
}
|
||||
@Test
|
||||
public void testTransferTableViaSSTables() throws Exception
|
||||
{
|
||||
doTransferTable(true);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -291,11 +322,11 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
state.writeElement(CounterId.fromInt(6), 3L, 3L);
|
||||
state.writeElement(CounterId.fromInt(8), 2L, 4L);
|
||||
cf.addColumn(new CounterColumn(ByteBufferUtil.bytes(col),
|
||||
state.context,
|
||||
timestamp));
|
||||
state.context,
|
||||
timestamp));
|
||||
cfCleaned.addColumn(new CounterColumn(ByteBufferUtil.bytes(col),
|
||||
cc.clearAllDelta(state.context),
|
||||
timestamp));
|
||||
cc.clearAllDelta(state.context),
|
||||
timestamp));
|
||||
|
||||
entries.put(key, cf);
|
||||
cleanedEntries.put(key, cfCleaned);
|
||||
|
|
@ -305,7 +336,7 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
.generation(0)
|
||||
.write(entries));
|
||||
}
|
||||
});
|
||||
}, true);
|
||||
|
||||
// filter pre-cleaned entries locally, and ensure that the end result is equal
|
||||
cleanedEntries.keySet().retainAll(keys);
|
||||
|
|
@ -319,7 +350,7 @@ public class StreamingTransferTest extends SchemaLoader
|
|||
|
||||
// Retransfer the file, making sure it is now idempotent (see CASSANDRA-3481)
|
||||
cfs.clearUnsafe();
|
||||
transfer(streamed);
|
||||
transferSSTables(streamed);
|
||||
SSTableReader restreamed = cfs.getSSTables().iterator().next();
|
||||
SSTableUtils.assertContentEquals(streamed, restreamed);
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue