mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-3.0' into cassandra-3.11
This commit is contained in:
commit
d6beb0113f
|
|
@ -3,6 +3,7 @@
|
|||
Merged from 3.0:
|
||||
=======
|
||||
3.0.21
|
||||
* Don't skip sstables in slice queries based only on local min/max/deletion timestamp (CASSANDRA-15690)
|
||||
* Memtable memory allocations may deadlock (CASSANDRA-15367)
|
||||
* Run evictFromMembership in GossipStage (CASSANDRA-15592)
|
||||
Merged from 2.2:
|
||||
|
|
|
|||
|
|
@ -720,18 +720,19 @@ public class SinglePartitionReadCommand extends ReadCommand
|
|||
* We can't eliminate full sstables based on the timestamp of what we've already read like
|
||||
* in collectTimeOrderedData, but we still want to eliminate sstable whose maxTimestamp < mostRecentTombstone
|
||||
* we've read. We still rely on the sstable ordering by maxTimestamp since if
|
||||
* maxTimestamp_s1 > maxTimestamp_s0,
|
||||
* maxTimestamp_s1 < maxTimestamp_s0,
|
||||
* we're guaranteed that s1 cannot have a row tombstone such that
|
||||
* timestamp(tombstone) > maxTimestamp_s0
|
||||
* since we necessarily have
|
||||
* timestamp(tombstone) <= maxTimestamp_s1
|
||||
* In other words, iterating in maxTimestamp order allow to do our mostRecentPartitionTombstone elimination
|
||||
* in one pass, and minimize the number of sstables for which we read a partition tombstone.
|
||||
*/
|
||||
* In other words, iterating in descending maxTimestamp order allow to do our mostRecentPartitionTombstone
|
||||
* elimination in one pass, and minimize the number of sstables for which we read a partition tombstone.
|
||||
*/
|
||||
Collections.sort(view.sstables, SSTableReader.maxTimestampDescending);
|
||||
long mostRecentPartitionTombstone = Long.MIN_VALUE;
|
||||
int nonIntersectingSSTables = 0;
|
||||
List<SSTableReader> skippedSSTablesWithTombstones = null;
|
||||
int includedDueToTombstones = 0;
|
||||
|
||||
SSTableReadMetricsCollector metricsCollector = new SSTableReadMetricsCollector();
|
||||
|
||||
for (SSTableReader sstable : view.sstables)
|
||||
|
|
@ -741,51 +742,47 @@ public class SinglePartitionReadCommand extends ReadCommand
|
|||
if (sstable.getMaxTimestamp() < mostRecentPartitionTombstone)
|
||||
break;
|
||||
|
||||
if (!shouldInclude(sstable))
|
||||
if (shouldInclude(sstable))
|
||||
{
|
||||
nonIntersectingSSTables++;
|
||||
if (sstable.mayHaveTombstones())
|
||||
{ // if sstable has tombstones we need to check after one pass if it can be safely skipped
|
||||
if (skippedSSTablesWithTombstones == null)
|
||||
skippedSSTablesWithTombstones = new ArrayList<>();
|
||||
skippedSSTablesWithTombstones.add(sstable);
|
||||
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
minTimestamp = Math.min(minTimestamp, sstable.getMinTimestamp());
|
||||
|
||||
@SuppressWarnings("resource") // 'iter' is added to iterators which is closed on exception,
|
||||
// or through the closing of the final merged iterator
|
||||
UnfilteredRowIteratorWithLowerBound iter = makeIterator(cfs, sstable, true, metricsCollector);
|
||||
if (!sstable.isRepaired())
|
||||
oldestUnrepairedTombstone = Math.min(oldestUnrepairedTombstone, sstable.getMinLocalDeletionTime());
|
||||
|
||||
iterators.add(iter);
|
||||
mostRecentPartitionTombstone = Math.max(mostRecentPartitionTombstone,
|
||||
iter.partitionLevelDeletion().markedForDeleteAt());
|
||||
}
|
||||
|
||||
int includedDueToTombstones = 0;
|
||||
// Check for sstables with tombstones that are not expired
|
||||
if (skippedSSTablesWithTombstones != null)
|
||||
{
|
||||
for (SSTableReader sstable : skippedSSTablesWithTombstones)
|
||||
{
|
||||
if (sstable.getMaxTimestamp() <= minTimestamp)
|
||||
continue;
|
||||
|
||||
@SuppressWarnings("resource") // 'iter' is added to iterators which is close on exception,
|
||||
// or through the closing of the final merged iterator
|
||||
UnfilteredRowIteratorWithLowerBound iter = makeIterator(cfs, sstable, false, metricsCollector);
|
||||
if (!sstable.isRepaired())
|
||||
oldestUnrepairedTombstone = Math.min(oldestUnrepairedTombstone, sstable.getMinLocalDeletionTime());
|
||||
|
||||
// 'iter' is added to iterators which is closed on exception, or through the closing of the final merged iterator
|
||||
@SuppressWarnings("resource")
|
||||
UnfilteredRowIterator iter = makeIterator(cfs, sstable, true, metricsCollector);
|
||||
iterators.add(iter);
|
||||
includedDueToTombstones++;
|
||||
mostRecentPartitionTombstone = Math.max(mostRecentPartitionTombstone,
|
||||
iter.partitionLevelDeletion().markedForDeleteAt());
|
||||
}
|
||||
else
|
||||
{
|
||||
|
||||
nonIntersectingSSTables++;
|
||||
// sstable contains no tombstone if maxLocalDeletionTime == Integer.MAX_VALUE, so we can safely skip those entirely
|
||||
if (sstable.mayHaveTombstones())
|
||||
{
|
||||
// 'iter' is added to iterators which is closed on exception, or through the closing of the final merged iterator
|
||||
@SuppressWarnings("resource")
|
||||
UnfilteredRowIterator iter = makeIterator(cfs, sstable, true, metricsCollector);
|
||||
// if the sstable contains a partition delete, then we must include it regardless of whether it
|
||||
// shadows any other data seen locally as we can't guarantee that other replicas have seen it
|
||||
if (!iter.partitionLevelDeletion().isLive())
|
||||
{
|
||||
if (!sstable.isRepaired())
|
||||
oldestUnrepairedTombstone = Math.min(oldestUnrepairedTombstone, sstable.getMinLocalDeletionTime());
|
||||
iterators.add(iter);
|
||||
includedDueToTombstones++;
|
||||
mostRecentPartitionTombstone = Math.max(mostRecentPartitionTombstone,
|
||||
iter.partitionLevelDeletion().markedForDeleteAt());
|
||||
}
|
||||
else
|
||||
{
|
||||
iter.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (Tracing.isTracing())
|
||||
Tracing.trace("Skipped {}/{} non-slice-intersecting sstables, included {} due to tombstones",
|
||||
nonIntersectingSSTables, view.sstables.size(), includedDueToTombstones);
|
||||
|
|
|
|||
|
|
@ -1,12 +1,16 @@
|
|||
package org.apache.cassandra.distributed.test;
|
||||
|
||||
import java.util.Set;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.apache.cassandra.db.Keyspace;
|
||||
import org.apache.cassandra.distributed.Cluster;
|
||||
import org.apache.cassandra.distributed.api.ConsistencyLevel;
|
||||
import org.apache.cassandra.distributed.api.ICluster;
|
||||
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
||||
import org.apache.cassandra.io.sstable.format.SSTableReader;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
|
||||
|
|
@ -269,6 +273,103 @@ public class SimpleReadWriteTest extends SharedClusterTestBase
|
|||
assertEquals(100, readCount1);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void skippedSSTableWithPartitionDeletionTest() throws Throwable
|
||||
{
|
||||
try (Cluster cluster = init(Cluster.create(2)))
|
||||
{
|
||||
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v int, PRIMARY KEY(pk, ck))");
|
||||
// insert a partition tombstone on node 1, the deletion timestamp should end up being the sstable's minTimestamp
|
||||
cluster.get(1).executeInternal("DELETE FROM " + KEYSPACE + ".tbl USING TIMESTAMP 1 WHERE pk = 0");
|
||||
// and a row from a different partition, to provide the sstable's min/max clustering
|
||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (1, 1, 1) USING TIMESTAMP 2");
|
||||
cluster.get(1).flush(KEYSPACE);
|
||||
// expect a single sstable, where minTimestamp equals the timestamp of the partition delete
|
||||
cluster.get(1).runOnInstance(() -> {
|
||||
Set<SSTableReader> sstables = Keyspace.open(KEYSPACE)
|
||||
.getColumnFamilyStore("tbl")
|
||||
.getLiveSSTables();
|
||||
assertEquals(1, sstables.size());
|
||||
assertEquals(1, sstables.iterator().next().getMinTimestamp());
|
||||
});
|
||||
|
||||
// on node 2, add a row for the deleted partition with an older timestamp than the deletion so it should be shadowed
|
||||
cluster.get(2).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (0, 10, 10) USING TIMESTAMP 0");
|
||||
|
||||
|
||||
Object[][] rows = cluster.coordinator(1)
|
||||
.execute("SELECT * FROM " + KEYSPACE + ".tbl WHERE pk=0 AND ck > 5",
|
||||
ConsistencyLevel.ALL);
|
||||
assertEquals(0, rows.length);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void skippedSSTableWithPartitionDeletionShadowingDataOnAnotherNode() throws Throwable
|
||||
{
|
||||
try (Cluster cluster = init(Cluster.create(2)))
|
||||
{
|
||||
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v int, PRIMARY KEY(pk, ck))");
|
||||
// insert a partition tombstone on node 1, the deletion timestamp should end up being the sstable's minTimestamp
|
||||
cluster.get(1).executeInternal("DELETE FROM " + KEYSPACE + ".tbl USING TIMESTAMP 1 WHERE pk = 0");
|
||||
// and a row from a different partition, to provide the sstable's min/max clustering
|
||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (1, 1, 1) USING TIMESTAMP 1");
|
||||
cluster.get(1).flush(KEYSPACE);
|
||||
// sstable 1 has minTimestamp == maxTimestamp == 1 and is skipped due to its min/max clusterings. Now we
|
||||
// insert a row which is not shadowed by the partition delete and flush to a second sstable. Importantly,
|
||||
// this sstable's minTimestamp is > than the maxTimestamp of the first sstable. This would cause the first
|
||||
// sstable not to be reincluded in the merge input, but we can't really make that decision as we don't
|
||||
// know what data and/or tombstones are present on other nodes
|
||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (0, 6, 6) USING TIMESTAMP 2");
|
||||
cluster.get(1).flush(KEYSPACE);
|
||||
|
||||
// on node 2, add a row for the deleted partition with an older timestamp than the deletion so it should be shadowed
|
||||
cluster.get(2).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (0, 10, 10) USING TIMESTAMP 0");
|
||||
|
||||
Object[][] rows = cluster.coordinator(1)
|
||||
.execute("SELECT * FROM " + KEYSPACE + ".tbl WHERE pk=0 AND ck > 5",
|
||||
ConsistencyLevel.ALL);
|
||||
// we expect that the row from node 2 (0, 10, 10) was shadowed by the partition delete, but the row from
|
||||
// node 1 (0, 6, 6) was not.
|
||||
assertRows(rows, new Object[] {0, 6 ,6});
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void skippedSSTableWithPartitionDeletionShadowingDataOnAnotherNode2() throws Throwable
|
||||
{
|
||||
// don't not add skipped sstables back just because the partition delete ts is < the local min ts
|
||||
|
||||
try (Cluster cluster = init(Cluster.create(2)))
|
||||
{
|
||||
cluster.schemaChange("CREATE TABLE " + KEYSPACE + ".tbl (pk int, ck int, v int, PRIMARY KEY(pk, ck))");
|
||||
// insert a partition tombstone on node 1, the deletion timestamp should end up being the sstable's minTimestamp
|
||||
cluster.get(1).executeInternal("DELETE FROM " + KEYSPACE + ".tbl USING TIMESTAMP 1 WHERE pk = 0");
|
||||
// and a row from a different partition, to provide the sstable's min/max clustering
|
||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (1, 1, 1) USING TIMESTAMP 3");
|
||||
cluster.get(1).flush(KEYSPACE);
|
||||
// sstable 1 has minTimestamp == maxTimestamp == 1 and is skipped due to its min/max clusterings. Now we
|
||||
// insert a row which is not shadowed by the partition delete and flush to a second sstable. The first sstable
|
||||
// has a maxTimestamp > than the min timestamp of all sstables, so it is a candidate for reinclusion to the
|
||||
// merge. Hoever, the second sstable's minTimestamp is > than the partition delete. This would cause the
|
||||
// first sstable not to be reincluded in the merge input, but we can't really make that decision as we don't
|
||||
// know what data and/or tombstones are present on other nodes
|
||||
cluster.get(1).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (0, 6, 6) USING TIMESTAMP 2");
|
||||
cluster.get(1).flush(KEYSPACE);
|
||||
|
||||
// on node 2, add a row for the deleted partition with an older timestamp than the deletion so it should be shadowed
|
||||
cluster.get(2).executeInternal("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v) VALUES (0, 10, 10) USING TIMESTAMP 0");
|
||||
|
||||
Object[][] rows = cluster.coordinator(1)
|
||||
.execute("SELECT * FROM " + KEYSPACE + ".tbl WHERE pk=0 AND ck > 5",
|
||||
ConsistencyLevel.ALL);
|
||||
// we expect that the row from node 2 (0, 10, 10) was shadowed by the partition delete, but the row from
|
||||
// node 1 (0, 6, 6) was not.
|
||||
assertRows(rows, new Object[] {0, 6 ,6});
|
||||
}
|
||||
}
|
||||
|
||||
private long readCount(IInvokableInstance instance)
|
||||
{
|
||||
return instance.callOnInstance(() -> Keyspace.open(KEYSPACE).getColumnFamilyStore("tbl").metric.readLatency.latency.getCount());
|
||||
|
|
|
|||
Loading…
Reference in New Issue