From 04e789bbcf283fe1d3b16613321ee169fa9f4e27 Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Mon, 7 Sep 2015 12:23:57 +0100 Subject: [PATCH 1/2] Cleanup, scrub and upgrade may unmark compacting early (CASSANDRA-10274) If an error occured during cleanup, scrub or upgrade (or any parallelAllSSTableOperation), the caller was immediately notified of the problem, and the method exited, executing the finally block that unmarked all of the sstables as compacting. Since the operations happen in parallel, many may still be running or waiting to run, and so another operation may operate over the same sstables, breaking the required mutual exclusivity. This patch ensures the method is not exited until all operations have completed, at which point the caller is notified of any exceptions. patch by benedict; reviewed by marcus for CASSANDRA-10274 --- CHANGES.txt | 2 ++ .../db/compaction/CompactionManager.java | 3 +-- .../apache/cassandra/utils/FBUtilities.java | 19 ++++++++++++++++--- .../apache/cassandra/utils/Throwables.java | 5 +++++ 4 files changed, 24 insertions(+), 5 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 6d43c9827e..e3ad5e84c7 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,6 @@ 2.1.10 + * Scrub, Cleanup and Upgrade do not unmark compacting until all operations + have completed, regardless of the occurence of exceptions (CASSANDRA-10274) * Fix handling of streaming EOF (CASSANDRA-10206) * Only check KeyCache when it is enabled * Change streaming_socket_timeout_in_ms default to 1 hour (CASSANDRA-8611) diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 3cfbe439c7..5d88a1174a 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -280,8 +280,7 @@ public class CompactionManager implements CompactionManagerMBean })); } - for (Future f : futures) - f.get(); + FBUtilities.waitOnFutures(futures); } finally { diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 1b118ba180..214f2f5b87 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -344,10 +344,23 @@ public class FBUtilities return System.currentTimeMillis() * 1000; } - public static void waitOnFutures(Iterable> futures) + public static List waitOnFutures(Iterable> futures) { - for (Future f : futures) - waitOnFuture(f); + List results = new ArrayList<>(); + Throwable fail = null; + for (Future f : futures) + { + try + { + results.add(f.get()); + } + catch (InterruptedException | ExecutionException e) + { + fail = Throwables.merge(fail, e); + } + } + Throwables.maybeFail(fail); + return results; } public static T waitOnFuture(Future future) diff --git a/src/java/org/apache/cassandra/utils/Throwables.java b/src/java/org/apache/cassandra/utils/Throwables.java index 552ca875b9..0a2bd2834a 100644 --- a/src/java/org/apache/cassandra/utils/Throwables.java +++ b/src/java/org/apache/cassandra/utils/Throwables.java @@ -29,4 +29,9 @@ public class Throwables return existingFail; } + public static void maybeFail(Throwable fail) + { + if (fail != null) + com.google.common.base.Throwables.propagate(fail); + } } From 21065423c6ce5a92194994bfa9e441d1d2b86b0a Mon Sep 17 00:00:00 2001 From: Brett Snyder Date: Wed, 9 Sep 2015 21:49:14 +0200 Subject: [PATCH 2/2] Fix update/delete behavior for static lists patch by Brett Snyder; reviewed by Benjamin Lerer for CASSANDRA-9838 --- .../statements/ModificationStatement.java | 6 ++++- .../validation/entities/CollectionsTest.java | 23 +++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java index 876c5e4ee6..37b46aeadc 100644 --- a/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java +++ b/src/java/org/apache/cassandra/cql3/statements/ModificationStatement.java @@ -29,6 +29,7 @@ import org.apache.cassandra.config.ColumnDefinition; import org.apache.cassandra.config.Schema; import org.apache.cassandra.cql3.*; import org.apache.cassandra.db.*; +import org.apache.cassandra.db.composites.AbstractCellNameType; import org.apache.cassandra.db.composites.CBuilder; import org.apache.cassandra.db.composites.Composite; import org.apache.cassandra.db.filter.ColumnSlice; @@ -446,7 +447,10 @@ public abstract class ModificationStatement implements CQLStatement if (row.cf == null || row.cf.isEmpty()) continue; - Iterator iter = cfm.comparator.CQL3RowBuilder(cfm, now).group(row.cf.getSortedColumns().iterator()); + CQL3Row.RowIterator iter = cfm.comparator.CQL3RowBuilder(cfm, now).group(row.cf.getSortedColumns().iterator()); + if(iter.getStaticRow() != null) { + map.put(row.key.getKey(), iter.getStaticRow()); + } if (iter.hasNext()) { map.put(row.key.getKey(), iter.next()); diff --git a/test/unit/org/apache/cassandra/cql3/validation/entities/CollectionsTest.java b/test/unit/org/apache/cassandra/cql3/validation/entities/CollectionsTest.java index 0241d4ffac..31dd5a6c00 100644 --- a/test/unit/org/apache/cassandra/cql3/validation/entities/CollectionsTest.java +++ b/test/unit/org/apache/cassandra/cql3/validation/entities/CollectionsTest.java @@ -485,4 +485,27 @@ public class CollectionsTest extends CQLTester assertInvalid("alter table %s add v set"); } + + /** + * Test for 9838. + */ + @Test + public void testUpdateStaticList() throws Throwable + { + createTable("CREATE TABLE %s (k1 text, k2 text, s_list list static, PRIMARY KEY (k1, k2))"); + + execute("insert into %s (k1, k2) VALUES ('a','b')"); + execute("update %s set s_list = s_list + [0] where k1='a'"); + assertRows(execute("select s_list from %s where k1='a'"), row(list(0))); + + execute("update %s set s_list[0] = 100 where k1='a'"); + assertRows(execute("select s_list from %s where k1='a'"), row(list(100))); + + execute("update %s set s_list = s_list + [0] where k1='a'"); + assertRows(execute("select s_list from %s where k1='a'"), row(list(100, 0))); + + execute("delete s_list[0] from %s where k1='a'"); + assertRows(execute("select s_list from %s where k1='a'"), row(list(0))); + } + }