From 195d6c76d87852a017e94f814c188bee7bb873fd Mon Sep 17 00:00:00 2001 From: finalchild Date: Thu, 19 Oct 2023 15:38:14 +0900 Subject: [PATCH 1/3] Remove duplicate paragraph in storage_engine.adoc --- .../cassandra/pages/architecture/storage_engine.adoc | 6 ------ 1 file changed, 6 deletions(-) diff --git a/doc/modules/cassandra/pages/architecture/storage_engine.adoc b/doc/modules/cassandra/pages/architecture/storage_engine.adoc index 9a0c37a089..47d1ff35b0 100644 --- a/doc/modules/cassandra/pages/architecture/storage_engine.adoc +++ b/doc/modules/cassandra/pages/architecture/storage_engine.adoc @@ -30,12 +30,6 @@ By default, max_mutation_size is half the size of `commitlog_segment_size`. `commitlog_segment_size` must be set to at least twice the size of `max_mutation_size`**. -Commitlogs are an append only log of all mutations local to a Cassandra -node. Any data written to Cassandra will first be written to a commit -log before being written to a memtable. This provides durability in the -case of unexpected shutdown. On startup, any mutations in the commit log -will be applied. - * `commitlog_sync`: may be either _periodic_ or _batch_. ** `batch`: In batch mode, Cassandra won’t ack writes until the commit log has been fsynced to disk. It will wait From 5e2468bf59581bcae9a55b4c3277dff2f8727a26 Mon Sep 17 00:00:00 2001 From: Josh McNeil <38532108+mcneiljt@users.noreply.github.com> Date: Wed, 19 Feb 2025 05:23:08 -0500 Subject: [PATCH 2/3] Fix broken linkage in Material View docs The `ALLOW FILTERING` link was missing a closing square bracket. --- doc/modules/cassandra/pages/cql/mvs.adoc | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/doc/modules/cassandra/pages/cql/mvs.adoc b/doc/modules/cassandra/pages/cql/mvs.adoc index 89ba0ceacb..fd0054a4c0 100644 --- a/doc/modules/cassandra/pages/cql/mvs.adoc +++ b/doc/modules/cassandra/pages/cql/mvs.adoc @@ -74,7 +74,7 @@ The `WHERE` clause has the following restrictions: ** no other restriction is allowed ** cannot have columns that are part of the _view_ primary key be null, they must always be at least restricted by a `IS NOT NULL` restriction (or any other restriction, but they must have one). -* cannot have an xref:cql/dml.adoc#ordering-clause[ordering clause], a xref:cql/dml.adoc#limit-clause[limit], or xref:cql/dml.adoc#allow-filtering[ALLOW FILTERING +* cannot have an xref:cql/dml.adoc#ordering-clause[ordering clause], a xref:cql/dml.adoc#limit-clause[limit], or xref:cql/dml.adoc#allow-filtering[ALLOW FILTERING] === MV primary key From 5aec56de3be35599746e232daf82f006c1addc16 Mon Sep 17 00:00:00 2001 From: Pedro Gordo Date: Fri, 10 Apr 2026 18:49:00 +0100 Subject: [PATCH 3/3] Backport CASSANDRA-17810 fix and improve RTBoundValidator error messages patch by Pedro Gordo; reviewed by Josh McKenzie, Stefan Miklosovic for CASSANDRA-18282 --- CHANGES.txt | 1 + .../org/apache/cassandra/db/ReadCommand.java | 82 +++++++++++-------- .../cassandra/db/ReadCommandVerbHandler.java | 26 ++++-- .../db/transform/RTBoundValidator.java | 50 ++++++----- .../exceptions/QueryCancelledException.java | 28 +++++++ .../cassandra/service/StorageProxy.java | 7 ++ .../distributed/test/TimeoutAbortTest.java | 62 ++++++++++++++ .../apache/cassandra/db/ReadCommandTest.java | 34 +++++++- 8 files changed, 227 insertions(+), 63 deletions(-) create mode 100644 src/java/org/apache/cassandra/exceptions/QueryCancelledException.java create mode 100644 test/distributed/org/apache/cassandra/distributed/test/TimeoutAbortTest.java diff --git a/CHANGES.txt b/CHANGES.txt index 31a9ce6c90..aa1466edd2 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0.21 + * Backport CASSANDRA-17810 fix and improve RTBoundValidator error messages (CASSANDRA-18282) 4.0.20 diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index bdbbb4cd74..9d3402767d 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -46,6 +46,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; @@ -391,7 +392,8 @@ public abstract class ReadCommand extends AbstractReadQuery try { - 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); @@ -557,9 +559,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) @@ -577,57 +579,71 @@ 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 void maybeDelayForTesting() - { - if (!metadata().keyspace.startsWith("system")) - FBUtilities.sleepQuietly(TEST_ITERATION_DELAY_MILLIS); } } - protected UnfilteredPartitionIterator withStateTracking(UnfilteredPartitionIterator iter) + private UnfilteredPartitionIterator withQueryCancellation(UnfilteredPartitionIterator iter) { - return Transformation.apply(iter, new CheckForAbort()); + 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); + } + + @Override + protected Row applyToRow(Row row) + { + FBUtilities.sleepQuietly(TEST_ITERATION_DELAY_MILLIS); + return row; + } + } + + private UnfilteredPartitionIterator maybeSlowDownForTesting(UnfilteredPartitionIterator iter) + { + 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 7bad9e6d11..3c5bb3c297 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; @@ -87,17 +90,28 @@ public class ReadCommandVerbHandler implements IVerbHandler { response = command.createResponse(iterator, controller.getRepairedDataInfo()); } + 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); + 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); - 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 8eec769fe9..f2d46042b3 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; @@ -229,7 +230,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 @@ -260,7 +270,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 @@ -291,7 +310,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