diff --git a/CHANGES.txt b/CHANGES.txt index eeb19d020f..1b7d7117d0 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,7 @@ 4.1.12 * Harden data resurrection startup check with atomic heartbeat file write with fallback (CASSANDRA-21290) +Merged from 4.0: + * Backport CASSANDRA-17810 fix and improve RTBoundValidator error messages (CASSANDRA-18282) 4.1.11 * Fix ant generate-eclipse-files (CASSANDRA-21215) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index b3a09e0249..5d90808ca2 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -48,6 +48,7 @@ import org.apache.cassandra.db.transform.RTBoundValidator; import org.apache.cassandra.db.transform.RTBoundValidator.Stage; import org.apache.cassandra.db.transform.StoppingTransformation; import org.apache.cassandra.db.transform.Transformation; +import org.apache.cassandra.exceptions.QueryCancelledException; import org.apache.cassandra.exceptions.UnknownIndexException; import org.apache.cassandra.index.Index; import org.apache.cassandra.index.IndexNotAvailableException; @@ -435,7 +436,8 @@ public abstract class ReadCommand extends AbstractReadQuery try { iterator = withQuerySizeTracking(iterator); - iterator = withStateTracking(iterator); + iterator = maybeSlowDownForTesting(iterator); + iterator = withQueryCancellation(iterator); iterator = RTBoundValidator.validate(withoutPurgeableTombstones(iterator, cfs, executionController), Stage.PURGED, false); iterator = withMetricsRecording(iterator, cfs.metric, startTimeNanos); @@ -614,9 +616,9 @@ public abstract class ReadCommand extends AbstractReadQuery return Transformation.apply(iter, new MetricRecording()); } - protected class CheckForAbort extends StoppingTransformation + private class QueryCancellationChecker extends StoppingTransformation { - long lastChecked = 0; + long lastCheckedAt = 0; @Override protected void attachTo(BasePartitions partitions) @@ -634,51 +636,62 @@ public abstract class ReadCommand extends AbstractReadQuery this.rows = rows; } + @Override protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition) { - if (maybeAbort()) - { - partition.close(); - return null; - } - + maybeCancel(); return Transformation.apply(partition, this); } + @Override protected Row applyToRow(Row row) { - if (TEST_ITERATION_DELAY_MILLIS > 0) - maybeDelayForTesting(); - - return maybeAbort() ? null : row; + maybeCancel(); + return row; } - private boolean maybeAbort() + private void maybeCancel() { - /** - * TODO: this is not a great way to abort early; why not expressly limit checks to 10ms intervals? + /* * The value returned by approxTime.now() is updated only every - * {@link org.apache.cassandra.utils.MonotonicClock.SampledClock.CHECK_INTERVAL_MS}, by default 2 millis. Since MonitorableImpl - * relies on approxTime, we don't need to check unless the approximate time has elapsed. + * {@link org.apache.cassandra.utils.MonotonicClock.SampledClock.CHECK_INTERVAL_MS}, by default 2 millis. + * Since MonitorableImpl relies on approxTime, we don't need to check unless the approximate time has elapsed. */ - if (lastChecked == approxTime.now()) - return false; - - lastChecked = approxTime.now(); + if (lastCheckedAt == approxTime.now()) + return; + lastCheckedAt = approxTime.now(); if (isAborted()) { stop(); - return true; + throw new QueryCancelledException(ReadCommand.this); } + } + } - return false; + private UnfilteredPartitionIterator withQueryCancellation(UnfilteredPartitionIterator iter) + { + return Transformation.apply(iter, new QueryCancellationChecker()); + } + + /** + * A transformation used for simulating slow queries by tests. + */ + @VisibleForTesting + private static class DelayInjector extends Transformation + { + @Override + protected UnfilteredRowIterator applyToPartition(UnfilteredRowIterator partition) + { + FBUtilities.sleepQuietly(TEST_ITERATION_DELAY_MILLIS); + return Transformation.apply(partition, this); } - private void maybeDelayForTesting() + @Override + protected Row applyToRow(Row row) { - if (!metadata().keyspace.startsWith("system")) - FBUtilities.sleepQuietly(TEST_ITERATION_DELAY_MILLIS); + FBUtilities.sleepQuietly(TEST_ITERATION_DELAY_MILLIS); + return row; } } @@ -766,9 +779,12 @@ public abstract class ReadCommand extends AbstractReadQuery return iterator; } - protected UnfilteredPartitionIterator withStateTracking(UnfilteredPartitionIterator iter) + private UnfilteredPartitionIterator maybeSlowDownForTesting(UnfilteredPartitionIterator iter) { - return Transformation.apply(iter, new CheckForAbort()); + if (TEST_ITERATION_DELAY_MILLIS > 0 && !SchemaConstants.isSystemKeyspace(metadata().keyspace)) + return Transformation.apply(iter, new DelayInjector()); + else + return iter; } /** diff --git a/src/java/org/apache/cassandra/db/ReadCommandVerbHandler.java b/src/java/org/apache/cassandra/db/ReadCommandVerbHandler.java index 7e28c8ab9e..726eba6dd3 100644 --- a/src/java/org/apache/cassandra/db/ReadCommandVerbHandler.java +++ b/src/java/org/apache/cassandra/db/ReadCommandVerbHandler.java @@ -19,6 +19,8 @@ package org.apache.cassandra.db; import java.util.concurrent.TimeUnit; +import com.google.common.base.Preconditions; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -26,6 +28,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.InvalidRequestException; +import org.apache.cassandra.exceptions.QueryCancelledException; import org.apache.cassandra.locator.Replica; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; @@ -106,18 +109,29 @@ public class ReadCommandVerbHandler implements IVerbHandler MessagingService.instance().send(reply, message.from()); return; } + catch (AssertionError t) + { + throw new AssertionError(String.format("Caught an error while trying to process the command: %s", command.toCQLString()), t); + } + catch (QueryCancelledException e) + { + logger.debug("Query cancelled (timeout)", e); + response = null; + Preconditions.checkState(!command.isCompleted(), "Read marked as completed despite being aborted by timeout to table %s", command.metadata()); + } - if (!command.complete()) + if (command.complete()) + { + Tracing.trace("Enqueuing response to {}", message.from()); + Message reply = message.responseWith(response); + reply = MessageParams.addToMessage(reply); + MessagingService.instance().send(reply, message.from()); + } + else { Tracing.trace("Discarding partial response to {} (timed out)", message.from()); MessagingService.instance().metrics.recordDroppedMessage(message, message.elapsedSinceCreated(NANOSECONDS), NANOSECONDS); - return; } - - Tracing.trace("Enqueuing response to {}", message.from()); - Message reply = message.responseWith(response); - reply = MessageParams.addToMessage(reply); - MessagingService.instance().send(reply, message.from()); } private void validateTransientStatus(Message message) diff --git a/src/java/org/apache/cassandra/db/transform/RTBoundValidator.java b/src/java/org/apache/cassandra/db/transform/RTBoundValidator.java index eb37f4bc1c..40fc826f3f 100644 --- a/src/java/org/apache/cassandra/db/transform/RTBoundValidator.java +++ b/src/java/org/apache/cassandra/db/transform/RTBoundValidator.java @@ -17,14 +17,16 @@ */ package org.apache.cassandra.db.transform; +import java.util.Objects; + import org.apache.cassandra.db.DeletionTime; import org.apache.cassandra.db.partitions.UnfilteredPartitionIterator; import org.apache.cassandra.db.rows.RangeTombstoneMarker; import org.apache.cassandra.db.rows.UnfilteredRowIterator; -import org.apache.cassandra.schema.TableMetadata; /** - * A validating transformation that sanity-checks the sequence of RT bounds and boundaries in every partition. + * A validating transformation that sanity-checks the sequence of Range Tombstone bounds and boundaries in every + * partition. * * What we validate, specifically: * - that open markers are only followed by close markers @@ -51,29 +53,27 @@ public final class RTBoundValidator extends Transformation errors = cluster.get(1).logs().grepForErrors().getResult(); + assertFalse(errors.toString(), errors.stream().anyMatch(s -> s.contains("open RT bound"))); + } + } +} diff --git a/test/unit/org/apache/cassandra/db/ReadCommandTest.java b/test/unit/org/apache/cassandra/db/ReadCommandTest.java index 43a7952175..bf272b87cf 100644 --- a/test/unit/org/apache/cassandra/db/ReadCommandTest.java +++ b/test/unit/org/apache/cassandra/db/ReadCommandTest.java @@ -56,6 +56,7 @@ import org.apache.cassandra.db.rows.UnfilteredRowIterators; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.exceptions.QueryCancelledException; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.io.util.DataInputBuffer; import org.apache.cassandra.io.util.DataOutputBuffer; @@ -232,7 +233,16 @@ public class ReadCommandTest assertEquals(2, Util.getAll(readCommand).size()); readCommand.abort(); - assertEquals(0, Util.getAll(readCommand).size()); + boolean cancelled = false; + try + { + Util.getAll(readCommand); + } + catch (QueryCancelledException e) + { + cancelled = true; + } + assertTrue(cancelled); } @Test @@ -263,7 +273,16 @@ public class ReadCommandTest assertEquals(2, partitions.get(0).rowCount()); readCommand.abort(); - assertEquals(0, Util.getAll(readCommand).size()); + boolean cancelled = false; + try + { + Util.getAll(readCommand); + } + catch (QueryCancelledException e) + { + cancelled = true; + } + assertTrue(cancelled); } @Test @@ -294,7 +313,16 @@ public class ReadCommandTest assertEquals(2, partitions.get(0).rowCount()); readCommand.abort(); - assertEquals(0, Util.getAll(readCommand).size()); + boolean cancelled = false; + try + { + Util.getAll(readCommand); + } + catch (QueryCancelledException e) + { + cancelled = true; + } + assertTrue(cancelled); } @Test