Merge branch 'cassandra-3.11' into trunk

* cassandra-3.11:
  Potential AssertionError during ReadRepair of range tombstone and partition deletions
This commit is contained in:
Sylvain Lebresne 2017-08-24 11:39:35 +02:00
commit 652d9f64f1
6 changed files with 193 additions and 16 deletions

View File

@ -129,6 +129,7 @@
* Duplicate the buffer before passing it to analyser in SASI operation (CASSANDRA-13512)
* Properly evict pstmts from prepared statements cache (CASSANDRA-13641)
Merged from 3.0:
* Potential AssertionError during ReadRepair of range tombstone and partition deletions (CASSANDRA-13719)
* Don't let stress write warmup data if n=0 (CASSANDRA-13773)
* Randomize batchlog endpoint selection with only 1 or 2 racks (CASSANDRA-12884)
* Fix digest calculation for counter cells (CASSANDRA-13750)

View File

@ -65,6 +65,28 @@ public abstract class ReadResponse
public abstract boolean isDigestResponse();
/**
* Creates a string of the requested partition in this read response suitable for debugging.
*/
public String toDebugString(ReadCommand command, DecoratedKey key)
{
if (isDigestResponse())
return "Digest:0x" + ByteBufferUtil.bytesToHex(digest(command));
try (UnfilteredPartitionIterator iter = makeIterator(command))
{
while (iter.hasNext())
{
try (UnfilteredRowIterator partition = iter.next())
{
if (partition.partitionKey().equals(key))
return ImmutableBTreePartition.create(partition).toString();
}
}
}
return "<key " + key + " not found>";
}
protected static ByteBuffer makeDigest(UnfilteredPartitionIterator iterator, ReadCommand command)
{
MessageDigest digest = FBUtilities.threadLocalMD5Digest();

View File

@ -94,7 +94,7 @@ public abstract class AbstractBTreePartition implements Partition, Iterable<Row>
public DeletionTime partitionLevelDeletion()
{
return holder().deletionInfo.getPartitionDeletion();
return deletionInfo().getPartitionDeletion();
}
public RegularAndStaticColumns columns()
@ -317,16 +317,20 @@ public abstract class AbstractBTreePartition implements Partition, Iterable<Row>
{
StringBuilder sb = new StringBuilder();
sb.append(String.format("[%s] key=%s columns=%s",
metadata().toString(),
sb.append(String.format("[%s] key=%s partition_deletion=%s columns=%s",
metadata(),
metadata().partitionKeyType.getString(partitionKey().getKey()),
partitionLevelDeletion(),
columns()));
if (staticRow() != Rows.EMPTY_STATIC_ROW)
sb.append("\n ").append(staticRow().toString(metadata()));
sb.append("\n ").append(staticRow().toString(metadata(), true));
for (Row row : this)
sb.append("\n ").append(row.toString(metadata()));
try (UnfilteredRowIterator iter = unfilteredIterator())
{
while (iter.hasNext())
sb.append("\n ").append(iter.next().toString(metadata(), true));
}
return sb.toString();
}

View File

@ -319,6 +319,15 @@ public class PartitionUpdate extends AbstractBTreePartition
return fromIterator(UnfilteredRowIterators.merge(asIterators, nowInSecs), ColumnFilter.all(updates.get(0).metadata()));
}
// We override this, because the version in the super-class calls holder(), which build the update preventing
// further updates, but that's not necessary here and being able to check at least the partition deletion without
// "locking" the update is nice (and used in DataResolver.RepairMergeListener.MergeListener).
@Override
public DeletionInfo deletionInfo()
{
return deletionInfo;
}
/**
* Modify this update to set every timestamp for live data to {@code newTimestamp} and
* every deletion timestamp to {@code newTimestamp - 1}.

View File

@ -22,6 +22,8 @@ import java.util.*;
import java.util.concurrent.TimeoutException;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Joiner;
import com.google.common.collect.Iterables;
import org.apache.cassandra.concurrent.Stage;
import org.apache.cassandra.concurrent.StageManager;
@ -245,6 +247,17 @@ public class DataResolver extends ResponseResolver
return repairs[i];
}
/**
* The partition level deletion with with which source {@code i} is currently repaired, or
* {@code DeletionTime.LIVE} if the source is not repaired on the partition level deletion (meaning it was
* up to date on it). The output* of this method is only valid after the call to
* {@link #onMergedPartitionLevelDeletion}.
*/
private DeletionTime partitionLevelRepairDeletion(int i)
{
return repairs[i] == null ? DeletionTime.LIVE : repairs[i].partitionLevelDeletion();
}
private Row.Builder currentRow(int i, Clustering clustering)
{
if (currentRows[i] == null)
@ -288,6 +301,37 @@ public class DataResolver extends ResponseResolver
}
public void onMergedRangeTombstoneMarkers(RangeTombstoneMarker merged, RangeTombstoneMarker[] versions)
{
try
{
// The code for merging range tombstones is a tad complex and we had the assertions there triggered
// unexpectedly in a few occasions (CASSANDRA-13237, CASSANDRA-13719). It's hard to get insights
// when that happen without more context that what the assertion errors give us however, hence the
// catch here that basically gather as much as context as reasonable.
internalOnMergedRangeTombstoneMarkers(merged, versions);
}
catch (AssertionError e)
{
// The following can be pretty verbose, but it's really only triggered if a bug happen, so we'd
// rather get more info to debug than not.
TableMetadata table = command.metadata();
String details = String.format("Error merging RTs on %s: merged=%s, versions=%s, sources={%s}, responses:%n %s",
table,
merged == null ? "null" : merged.toString(table),
'[' + Joiner.on(", ").join(Iterables.transform(Arrays.asList(versions), rt -> rt == null ? "null" : rt.toString(table))) + ']',
Arrays.toString(sources),
makeResponsesDebugString());
throw new AssertionError(details, e);
}
}
private String makeResponsesDebugString()
{
return Joiner.on(",\n")
.join(Iterables.transform(getMessages(), m -> m.from + " => " + m.payload.toDebugString(command, partitionKey)));
}
private void internalOnMergedRangeTombstoneMarkers(RangeTombstoneMarker merged, RangeTombstoneMarker[] versions)
{
// The current deletion as of dealing with this marker.
DeletionTime currentDeletion = currentDeletion();
@ -313,21 +357,27 @@ public class DataResolver extends ResponseResolver
// active after that point. Further whatever deletion was open or is open by this marker on the
// source, that deletion cannot supersedes the current one.
//
// But while the marker deletion (before and/or after this point) cannot supersed the current
// But while the marker deletion (before and/or after this point) cannot supersede the current
// deletion, we want to know if it's equal to it (both before and after), because in that case
// the source is up to date and we don't want to include repair.
//
// So in practice we have 2 possible case:
// 1) the source was up-to-date on deletion up to that point (markerToRepair[i] == null). Then
// it won't be from that point on unless it's a boundary and the new opened deletion time
// is also equal to the current deletion (note that this implies the boundary has the same
// closing and opening deletion time, which should generally not happen, but can due to legacy
// reading code not avoiding this for a while, see CASSANDRA-13237).
// 2) the source wasn't up-to-date on deletion up to that point (markerToRepair[i] != null), and
// it may now be (if it isn't we just have nothing to do for that marker).
// 1) the source was up-to-date on deletion up to that point: then it won't be from that point
// on unless it's a boundary and the new opened deletion time is also equal to the current
// deletion (note that this implies the boundary has the same closing and opening deletion
// time, which should generally not happen, but can due to legacy reading code not avoiding
// this for a while, see CASSANDRA-13237).
// 2) the source wasn't up-to-date on deletion up to that point and it may now be (if it isn't
// we just have nothing to do for that marker).
assert !currentDeletion.isLive() : currentDeletion.toString();
if (markerToRepair[i] == null)
// Is the source up to date on deletion? It's up to date if it doesn't have an open RT repair
// nor an "active" partition level deletion (where "active" means that it's greater or equal
// to the current deletion: if the source has a repaired partition deletion lower than the
// current deletion, this means the current deletion is due to a previously open range tombstone,
// and if the source isn't currently repaired for that RT, then it means it's up to date on it).
DeletionTime partitionRepairDeletion = partitionLevelRepairDeletion(i);
if (markerToRepair[i] == null && currentDeletion.supersedes(partitionRepairDeletion))
{
// Since there is an ongoing merged deletion, the only way we don't have an open repair for
// this source is that it had a range open with the same deletion as current and it's
@ -342,6 +392,8 @@ public class DataResolver extends ResponseResolver
markerToRepair[i] = marker.closeBound(isReversed).invert();
}
// In case 2) above, we only have something to do if the source is up-to-date after that point
// (which, since the source isn't up-to-date before that point, means we're opening a new deletion
// that is equal to the current one).
else if (marker.isOpen(isReversed) && currentDeletion.equals(marker.openDeletionTime(isReversed)))
{
closeOpenMarker(i, marker.openBound(isReversed).invert());

View File

@ -576,7 +576,7 @@ public class DataResolverTest
* same deletion on both side (while is useless but could be created by legacy code pre-CASSANDRA-13237 and could
* thus still be sent).
*/
public void testRepairRangeTombstoneBoundary(int timestamp1, int timestamp2, int timestamp3) throws UnknownHostException
private void testRepairRangeTombstoneBoundary(int timestamp1, int timestamp2, int timestamp3) throws UnknownHostException
{
DataResolver resolver = new DataResolver(ks, command, ConsistencyLevel.ALL, 2, System.nanoTime());
InetAddress peer1 = peer();
@ -623,6 +623,95 @@ public class DataResolverTest
assertRepairContainsDeletions(msg, null, expected);
}
/**
* Test for CASSANDRA-13719: tests that having a partition deletion shadow a range tombstone on another source
* doesn't trigger an assertion error.
*/
@Test
public void testRepairRangeTombstoneWithPartitionDeletion()
{
DataResolver resolver = new DataResolver(ks, command, ConsistencyLevel.ALL, 2, System.nanoTime());
InetAddress peer1 = peer();
InetAddress peer2 = peer();
// 1st "stream": just a partition deletion
UnfilteredPartitionIterator iter1 = iter(PartitionUpdate.fullPartitionDelete(cfm, dk, 10, nowInSec));
// 2nd "stream": a range tombstone that is covered by the 1st stream
RangeTombstone rt = tombstone("0", true , "10", true, 5, nowInSec);
UnfilteredPartitionIterator iter2 = iter(new RowUpdateBuilder(cfm, nowInSec, 1L, dk)
.addRangeTombstone(rt)
.buildUpdate());
resolver.preprocess(readResponseMessage(peer1, iter1));
resolver.preprocess(readResponseMessage(peer2, iter2));
// No results, we've only reconciled tombstones.
try (PartitionIterator data = resolver.resolve())
{
assertFalse(data.hasNext());
// 2nd stream should get repaired
assertRepairFuture(resolver, 1);
}
assertEquals(1, messageRecorder.sent.size());
MessageOut msg = getSentMessage(peer2);
assertRepairMetadata(msg);
assertRepairContainsNoColumns(msg);
assertRepairContainsDeletions(msg, new DeletionTime(10, nowInSec));
}
/**
* Additional test for CASSANDRA-13719: tests the case where a partition deletion doesn't shadow a range tombstone.
*/
@Test
public void testRepairRangeTombstoneWithPartitionDeletion2()
{
DataResolver resolver = new DataResolver(ks, command, ConsistencyLevel.ALL, 2, System.nanoTime());
InetAddress peer1 = peer();
InetAddress peer2 = peer();
// 1st "stream": a partition deletion and a range tombstone
RangeTombstone rt1 = tombstone("0", true , "9", true, 11, nowInSec);
PartitionUpdate upd1 = new RowUpdateBuilder(cfm, nowInSec, 1L, dk)
.addRangeTombstone(rt1)
.buildUpdate();
((MutableDeletionInfo)upd1.deletionInfo()).add(new DeletionTime(10, nowInSec));
UnfilteredPartitionIterator iter1 = iter(upd1);
// 2nd "stream": a range tombstone that is covered by the other stream rt
RangeTombstone rt2 = tombstone("2", true , "3", true, 11, nowInSec);
RangeTombstone rt3 = tombstone("4", true , "5", true, 10, nowInSec);
UnfilteredPartitionIterator iter2 = iter(new RowUpdateBuilder(cfm, nowInSec, 1L, dk)
.addRangeTombstone(rt2)
.addRangeTombstone(rt3)
.buildUpdate());
resolver.preprocess(readResponseMessage(peer1, iter1));
resolver.preprocess(readResponseMessage(peer2, iter2));
// No results, we've only reconciled tombstones.
try (PartitionIterator data = resolver.resolve())
{
assertFalse(data.hasNext());
// 2nd stream should get repaired
assertRepairFuture(resolver, 1);
}
assertEquals(1, messageRecorder.sent.size());
MessageOut msg = getSentMessage(peer2);
assertRepairMetadata(msg);
assertRepairContainsNoColumns(msg);
// 2nd stream should get both the partition deletion, as well as the part of the 1st stream RT that it misses
assertRepairContainsDeletions(msg, new DeletionTime(10, nowInSec),
tombstone("0", true, "2", false, 11, nowInSec),
tombstone("3", false, "9", true, 11, nowInSec));
}
// Forces the start to be exclusive if the condition holds
private static RangeTombstone withExclusiveStartIf(RangeTombstone rt, boolean condition)
{