From f8452838a964c7bfa938277c9fdfc337f9aff886 Mon Sep 17 00:00:00 2001 From: Marcus Eriksson Date: Mon, 14 Dec 2015 15:33:21 +0100 Subject: [PATCH 01/16] alter/drop user should be case sensitive Patch by marcuse; reviewed by Sam Tunnicliffe for CASSANDRA-10817 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/cql3/Cql.g | 4 ++-- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index b700102a90..bb5909c06e 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.2.5 + * drop/alter user should be case sensitive (CASSANDRA-10817) * jemalloc detection fails due to quoting issues in regexv (CASSANDRA-10946) * Support counter-columns for native aggregates (sum,avg,max,min) (CASSANDRA-9977) * (cqlsh) show correct column names for empty result sets (CASSANDRA-9813) diff --git a/src/java/org/apache/cassandra/cql3/Cql.g b/src/java/org/apache/cassandra/cql3/Cql.g index b4cbac805e..035e7044d9 100644 --- a/src/java/org/apache/cassandra/cql3/Cql.g +++ b/src/java/org/apache/cassandra/cql3/Cql.g @@ -982,7 +982,7 @@ alterUserStatement returns [AlterRoleStatement stmt] RoleOptions opts = new RoleOptions(); RoleName name = new RoleName(); } - : K_ALTER K_USER u=username { name.setName($u.text, false); } + : K_ALTER K_USER u=username { name.setName($u.text, true); } ( K_WITH userPassword[opts] )? ( K_SUPERUSER { opts.setOption(IRoleManager.Option.SUPERUSER, true); } | K_NOSUPERUSER { opts.setOption(IRoleManager.Option.SUPERUSER, false); } ) ? @@ -997,7 +997,7 @@ dropUserStatement returns [DropRoleStatement stmt] boolean ifExists = false; RoleName name = new RoleName(); } - : K_DROP K_USER (K_IF K_EXISTS { ifExists = true; })? u=username { name.setName($u.text, false); $stmt = new DropRoleStatement(name, ifExists); } + : K_DROP K_USER (K_IF K_EXISTS { ifExists = true; })? u=username { name.setName($u.text, true); $stmt = new DropRoleStatement(name, ifExists); } ; /** From 9dafa438a5dcce8674aaa945b1495e70f95d2839 Mon Sep 17 00:00:00 2001 From: Carl Yeksigian Date: Fri, 18 Dec 2015 10:59:44 -0500 Subject: [PATCH 02/16] Skip commit log and saved cache directories in SSTable version startup check Patch by Carl Yeksigian; reviewed by Sam Tunnicliffe for CASSANDRA-10902 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/service/StartupChecks.java | 9 ++++++++- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index bb5909c06e..3c919c7654 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.2.5 + * Skip commit log and saved cache directories in SSTable version startup check (CASSANDRA-10902) * drop/alter user should be case sensitive (CASSANDRA-10817) * jemalloc detection fails due to quoting issues in regexv (CASSANDRA-10946) * Support counter-columns for native aggregates (sum,avg,max,min) (CASSANDRA-9977) diff --git a/src/java/org/apache/cassandra/service/StartupChecks.java b/src/java/org/apache/cassandra/service/StartupChecks.java index c13b401e08..19f32b6dc7 100644 --- a/src/java/org/apache/cassandra/service/StartupChecks.java +++ b/src/java/org/apache/cassandra/service/StartupChecks.java @@ -36,6 +36,7 @@ import org.apache.cassandra.db.*; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.StartupException; import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.utils.*; /** @@ -233,6 +234,10 @@ public class StartupChecks public void execute() throws StartupException { final Set invalid = new HashSet<>(); + final Set nonSSTablePaths = new HashSet<>(); + nonSSTablePaths.add(FileUtils.getCanonicalPath(DatabaseDescriptor.getCommitLogLocation())); + nonSSTablePaths.add(FileUtils.getCanonicalPath(DatabaseDescriptor.getSavedCachesLocation())); + FileVisitor sstableVisitor = new SimpleFileVisitor() { public FileVisitResult visitFile(Path file, BasicFileAttributes attrs) throws IOException @@ -255,7 +260,9 @@ public class StartupChecks public FileVisitResult preVisitDirectory(Path dir, BasicFileAttributes attrs) throws IOException { String name = dir.getFileName().toString(); - return (name.equals("snapshots") || name.equals("backups")) + return (name.equals("snapshots") + || name.equals("backups") + || nonSSTablePaths.contains(dir.toFile().getCanonicalPath())) ? FileVisitResult.SKIP_SUBTREE : FileVisitResult.CONTINUE; } From 1d7bacc45fa1cd6cac36d7f9ece30ba1ed430f2a Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Mon, 28 Dec 2015 14:08:10 +0100 Subject: [PATCH 03/16] Fix potential assertion error during compaction patch by slebresne; reviewed by krummas for CASSANDRA-10944 --- CHANGES.txt | 1 + .../org/apache/cassandra/db/ReadCommand.java | 2 +- .../db/compaction/CompactionIterator.java | 6 +-- .../db/partitions/PurgeFunction.java | 12 ++--- .../apache/cassandra/db/rows/BufferCell.java | 2 +- .../db/compaction/TTLExpiryTest.java | 47 ++++++++++++++----- 6 files changed, 46 insertions(+), 24 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index beb59b0936..23edbbf329 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 3.0.3 + * Fix potential assertion error during compaction (CASSANDRA-10944) * Fix counting of received sstables in streaming (CASSANDRA-10949) * Implement hints compression (CASSANDRA-9428) * Fix potential assertion error when reading static columns (CASSANDRA-10903) diff --git a/src/java/org/apache/cassandra/db/ReadCommand.java b/src/java/org/apache/cassandra/db/ReadCommand.java index 5ab1ee565a..3f0695c79d 100644 --- a/src/java/org/apache/cassandra/db/ReadCommand.java +++ b/src/java/org/apache/cassandra/db/ReadCommand.java @@ -488,7 +488,7 @@ public abstract class ReadCommand implements ReadQuery { public WithoutPurgeableTombstones() { - super(isForThrift, cfs.gcBefore(nowInSec()), oldestUnrepairedTombstone(), cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones()); + super(isForThrift, nowInSec(), cfs.gcBefore(nowInSec()), oldestUnrepairedTombstone(), cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones()); } protected long getMaxPurgeableTimestamp() diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index 8a3b24b2f3..d39da2afce 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -103,7 +103,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte ? EmptyIterators.unfilteredPartition(controller.cfs.metadata, false) : UnfilteredPartitionIterators.merge(scanners, nowInSec, listener()); boolean isForThrift = merged.isForThrift(); // to stop capture of iterator in Purger, which is confusing for debug - this.compacted = Transformation.apply(merged, new Purger(isForThrift, controller)); + this.compacted = Transformation.apply(merged, new Purger(isForThrift, controller, nowInSec)); } public boolean isForThrift() @@ -264,9 +264,9 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte private long compactedUnfiltered; - private Purger(boolean isForThrift, CompactionController controller) + private Purger(boolean isForThrift, CompactionController controller, int nowInSec) { - super(isForThrift, controller.gcBefore, controller.compactingRepaired() ? Integer.MIN_VALUE : Integer.MAX_VALUE, controller.cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones()); + super(isForThrift, nowInSec, controller.gcBefore, controller.compactingRepaired() ? Integer.MIN_VALUE : Integer.MAX_VALUE, controller.cfs.getCompactionStrategyManager().onlyPurgeRepairedTombstones()); this.controller = controller; } diff --git a/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java b/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java index b7b01d6490..492bab1c7d 100644 --- a/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java +++ b/src/java/org/apache/cassandra/db/partitions/PurgeFunction.java @@ -25,13 +25,13 @@ public abstract class PurgeFunction extends Transformation !(onlyPurgeRepairedTombstones && localDeletionTime >= oldestUnrepairedTombstone) && localDeletionTime < gcBefore @@ -79,13 +79,13 @@ public abstract class PurgeFunction extends Transformation Date: Tue, 5 Jan 2016 15:22:06 +0100 Subject: [PATCH 04/16] Optimize pending ranges computation patch by dikanggu; reviewed by blambov for CASSANDRA-9258 --- .../cassandra/locator/TokenMetadata.java | 61 ++++++++++--------- .../cassandra/service/StorageService.java | 2 +- 2 files changed, 32 insertions(+), 31 deletions(-) diff --git a/src/java/org/apache/cassandra/locator/TokenMetadata.java b/src/java/org/apache/cassandra/locator/TokenMetadata.java index db0b609ffa..00d8ee978d 100644 --- a/src/java/org/apache/cassandra/locator/TokenMetadata.java +++ b/src/java/org/apache/cassandra/locator/TokenMetadata.java @@ -82,7 +82,7 @@ public class TokenMetadata // (don't need to record Token here since it's still part of tokenToEndpointMap until it's done leaving) private final Set leavingEndpoints = new HashSet<>(); // this is a cache of the calculation from {tokenToEndpointMap, bootstrapTokens, leavingEndpoints} - private final ConcurrentMap, InetAddress>> pendingRanges = new ConcurrentHashMap<>(); + private final ConcurrentMap pendingRanges = new ConcurrentHashMap(); // nodes which are migrating to the new tokens in the ring private final Set> movingEndpoints = new HashSet<>(); @@ -673,23 +673,30 @@ public class TokenMetadata return sortedTokens; } - private Multimap, InetAddress> getPendingRangesMM(String keyspaceName) + public Multimap, InetAddress> getPendingRangesMM(String keyspaceName) { - Multimap, InetAddress> map = pendingRanges.get(keyspaceName); - if (map == null) + Multimap, InetAddress> map = HashMultimap.create(); + PendingRangeMaps pendingRangeMaps = this.pendingRanges.get(keyspaceName); + + if (pendingRangeMaps != null) { - map = HashMultimap.create(); - Multimap, InetAddress> priorMap = pendingRanges.putIfAbsent(keyspaceName, map); - if (priorMap != null) - map = priorMap; + for (Map.Entry, List> entry : pendingRangeMaps) + { + Range range = entry.getKey(); + for (InetAddress address : entry.getValue()) + { + map.put(range, address); + } + } } + return map; } /** a mutable map may be returned but caller should not modify it */ - public Map, Collection> getPendingRanges(String keyspaceName) + public PendingRangeMaps getPendingRanges(String keyspaceName) { - return getPendingRangesMM(keyspaceName).asMap(); + return this.pendingRanges.get(keyspaceName); } public List> getPendingRanges(String keyspaceName, InetAddress endpoint) @@ -733,7 +740,7 @@ public class TokenMetadata lock.readLock().lock(); try { - Multimap, InetAddress> newPendingRanges = HashMultimap.create(); + PendingRangeMaps newPendingRanges = new PendingRangeMaps(); if (bootstrapTokens.isEmpty() && leavingEndpoints.isEmpty() && movingEndpoints.isEmpty()) { @@ -761,7 +768,10 @@ public class TokenMetadata { Set currentEndpoints = ImmutableSet.copyOf(strategy.calculateNaturalEndpoints(range.right, metadata)); Set newEndpoints = ImmutableSet.copyOf(strategy.calculateNaturalEndpoints(range.right, allLeftMetadata)); - newPendingRanges.putAll(range, Sets.difference(newEndpoints, currentEndpoints)); + for (InetAddress address : Sets.difference(newEndpoints, currentEndpoints)) + { + newPendingRanges.addPendingRange(range, address); + } } // At this stage newPendingRanges has been updated according to leave operations. We can @@ -776,7 +786,9 @@ public class TokenMetadata allLeftMetadata.updateNormalTokens(tokens, endpoint); for (Range range : strategy.getAddressRanges(allLeftMetadata).get(endpoint)) - newPendingRanges.put(range, endpoint); + { + newPendingRanges.addPendingRange(range, endpoint); + } allLeftMetadata.removeEndpoint(endpoint); } @@ -794,7 +806,7 @@ public class TokenMetadata for (Range range : strategy.getAddressRanges(allLeftMetadata).get(endpoint)) { - newPendingRanges.put(range, endpoint); + newPendingRanges.addPendingRange(range, endpoint); } allLeftMetadata.removeEndpoint(endpoint); @@ -1029,13 +1041,9 @@ public class TokenMetadata { StringBuilder sb = new StringBuilder(); - for (Map.Entry, InetAddress>> entry : pendingRanges.entrySet()) + for (PendingRangeMaps pendingRangeMaps : pendingRanges.values()) { - for (Map.Entry, InetAddress> rmap : entry.getValue().entries()) - { - sb.append(rmap.getValue()).append(':').append(rmap.getKey()); - sb.append(System.getProperty("line.separator")); - } + sb.append(pendingRangeMaps.printPendingRanges()); } return sb.toString(); @@ -1043,18 +1051,11 @@ public class TokenMetadata public Collection pendingEndpointsFor(Token token, String keyspaceName) { - Map, Collection> ranges = getPendingRanges(keyspaceName); - if (ranges.isEmpty()) + PendingRangeMaps pendingRangeMaps = this.pendingRanges.get(keyspaceName); + if (pendingRangeMaps == null) return Collections.emptyList(); - Set endpoints = new HashSet<>(); - for (Map.Entry, Collection> entry : ranges.entrySet()) - { - if (entry.getKey().contains(token)) - endpoints.addAll(entry.getValue()); - } - - return endpoints; + return pendingRangeMaps.pendingEndpointsFor(token); } /** diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index e8e7daf358..84ebd9a82a 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -1378,7 +1378,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE keyspace = Schema.instance.getNonSystemKeyspaces().get(0); Map, List> map = new HashMap<>(); - for (Map.Entry, Collection> entry : tokenMetadata.getPendingRanges(keyspace).entrySet()) + for (Map.Entry, Collection> entry : tokenMetadata.getPendingRangesMM(keyspace).asMap().entrySet()) { List l = new ArrayList<>(entry.getValue()); map.put(entry.getKey().asList(), stringify(l)); From e0c1b0bb7121df1cc0185ffc0b35547f75daa281 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Tue, 5 Jan 2016 15:26:54 +0100 Subject: [PATCH 05/16] Add mistakenly forgotten files for CASSANDRA-9258 --- CHANGES.txt | 1 + .../cassandra/locator/PendingRangeMaps.java | 209 ++++++++++++++++++ .../test/microbench/PendingRangesBench.java | 89 ++++++++ .../locator/PendingRangeMapsTest.java | 78 +++++++ 4 files changed, 377 insertions(+) create mode 100644 src/java/org/apache/cassandra/locator/PendingRangeMaps.java create mode 100644 test/microbench/org/apache/cassandra/test/microbench/PendingRangesBench.java create mode 100644 test/unit/org/apache/cassandra/locator/PendingRangeMapsTest.java diff --git a/CHANGES.txt b/CHANGES.txt index 3c919c7654..648200bbca 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.2.5 + * Optimize pending range computation (CASSANDRA-9258) * Skip commit log and saved cache directories in SSTable version startup check (CASSANDRA-10902) * drop/alter user should be case sensitive (CASSANDRA-10817) * jemalloc detection fails due to quoting issues in regexv (CASSANDRA-10946) diff --git a/src/java/org/apache/cassandra/locator/PendingRangeMaps.java b/src/java/org/apache/cassandra/locator/PendingRangeMaps.java new file mode 100644 index 0000000000..1892cc376a --- /dev/null +++ b/src/java/org/apache/cassandra/locator/PendingRangeMaps.java @@ -0,0 +1,209 @@ +package org.apache.cassandra.locator; + +import com.google.common.collect.Iterators; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.net.InetAddress; +import java.util.*; + +public class PendingRangeMaps implements Iterable, List>> +{ + private static final Logger logger = LoggerFactory.getLogger(PendingRangeMaps.class); + + /** + * We have for NavigableMap to be able to search for ranges containing a token efficiently. + * + * First two are for non-wrap-around ranges, and the last two are for wrap-around ranges. + */ + // ascendingMap will sort the ranges by the ascending order of right token + final NavigableMap, List> ascendingMap; + /** + * sorting end ascending, if ends are same, sorting begin descending, so that token (end, end) will + * come before (begin, end] with the same end, and (begin, end) will be selected in the tailMap. + */ + static final Comparator> ascendingComparator = new Comparator>() + { + @Override + public int compare(Range o1, Range o2) + { + int res = o1.right.compareTo(o2.right); + if (res != 0) + return res; + + return o2.left.compareTo(o1.left); + } + }; + + // ascendingMap will sort the ranges by the descending order of left token + final NavigableMap, List> descendingMap; + /** + * sorting begin descending, if begins are same, sorting end descending, so that token (begin, begin) will + * come after (begin, end] with the same begin, and (begin, end) won't be selected in the tailMap. + */ + static final Comparator> descendingComparator = new Comparator>() + { + @Override + public int compare(Range o1, Range o2) + { + int res = o2.left.compareTo(o1.left); + if (res != 0) + return res; + + // if left tokens are same, sort by the descending of the right tokens. + return o2.right.compareTo(o1.right); + } + }; + + // these two maps are for warp around ranges. + final NavigableMap, List> ascendingMapForWrapAround; + /** + * for wrap around range (begin, end], which begin > end. + * Sorting end ascending, if ends are same, sorting begin ascending, + * so that token (end, end) will come before (begin, end] with the same end, and (begin, end] will be selected in + * the tailMap. + */ + static final Comparator> ascendingComparatorForWrapAround = new Comparator>() + { + @Override + public int compare(Range o1, Range o2) + { + int res = o1.right.compareTo(o2.right); + if (res != 0) + return res; + + return o1.left.compareTo(o2.left); + } + }; + + final NavigableMap, List> descendingMapForWrapAround; + /** + * for wrap around ranges, which begin > end. + * Sorting end ascending, so that token (begin, begin) will come after (begin, end] with the same begin, + * and (begin, end) won't be selected in the tailMap. + */ + static final Comparator> descendingComparatorForWrapAround = new Comparator>() + { + @Override + public int compare(Range o1, Range o2) + { + int res = o2.left.compareTo(o1.left); + if (res != 0) + return res; + return o1.right.compareTo(o2.right); + } + }; + + public PendingRangeMaps() + { + this.ascendingMap = new TreeMap, List>(ascendingComparator); + this.descendingMap = new TreeMap, List>(descendingComparator); + this.ascendingMapForWrapAround = new TreeMap, List>(ascendingComparatorForWrapAround); + this.descendingMapForWrapAround = new TreeMap, List>(descendingComparatorForWrapAround); + } + + static final void addToMap(Range range, + InetAddress address, + NavigableMap, List> ascendingMap, + NavigableMap, List> descendingMap) + { + List addresses = ascendingMap.get(range); + if (addresses == null) + { + addresses = new ArrayList(1); + ascendingMap.put(range, addresses); + descendingMap.put(range, addresses); + } + addresses.add(address); + } + + public void addPendingRange(Range range, InetAddress address) + { + if (Range.isWrapAround(range.left, range.right)) + { + addToMap(range, address, ascendingMapForWrapAround, descendingMapForWrapAround); + } + else + { + addToMap(range, address, ascendingMap, descendingMap); + } + } + + static final void addIntersections(Set endpointsToAdd, + NavigableMap, List> smallerMap, + NavigableMap, List> biggerMap) + { + // find the intersection of two sets + for (Range range : smallerMap.keySet()) + { + List addresses = biggerMap.get(range); + if (addresses != null) + { + endpointsToAdd.addAll(addresses); + } + } + } + + public Collection pendingEndpointsFor(Token token) + { + Set endpoints = new HashSet<>(); + + Range searchRange = new Range(token, token); + + // search for non-wrap-around maps + NavigableMap, List> ascendingTailMap = ascendingMap.tailMap(searchRange, true); + NavigableMap, List> descendingTailMap = descendingMap.tailMap(searchRange, false); + + // add intersections of two maps + if (ascendingTailMap.size() < descendingTailMap.size()) + { + addIntersections(endpoints, ascendingTailMap, descendingTailMap); + } + else + { + addIntersections(endpoints, descendingTailMap, ascendingTailMap); + } + + // search for wrap-around sets + ascendingTailMap = ascendingMapForWrapAround.tailMap(searchRange, true); + descendingTailMap = descendingMapForWrapAround.tailMap(searchRange, false); + + // add them since they are all necessary. + for (Map.Entry, List> entry : ascendingTailMap.entrySet()) + { + endpoints.addAll(entry.getValue()); + } + for (Map.Entry, List> entry : descendingTailMap.entrySet()) + { + endpoints.addAll(entry.getValue()); + } + + return endpoints; + } + + public String printPendingRanges() + { + StringBuilder sb = new StringBuilder(); + + for (Map.Entry, List> entry : this) + { + Range range = entry.getKey(); + + for (InetAddress address : entry.getValue()) + { + sb.append(address).append(':').append(range); + sb.append(System.getProperty("line.separator")); + } + } + + return sb.toString(); + } + + @Override + public Iterator, List>> iterator() + { + return Iterators.concat(ascendingMap.entrySet().iterator(), ascendingMapForWrapAround.entrySet().iterator()); + } +} diff --git a/test/microbench/org/apache/cassandra/test/microbench/PendingRangesBench.java b/test/microbench/org/apache/cassandra/test/microbench/PendingRangesBench.java new file mode 100644 index 0000000000..e50cbafcfa --- /dev/null +++ b/test/microbench/org/apache/cassandra/test/microbench/PendingRangesBench.java @@ -0,0 +1,89 @@ +package org.apache.cassandra.test.microbench; + +import com.google.common.collect.HashMultimap; +import com.google.common.collect.Multimap; +import org.apache.cassandra.dht.RandomPartitioner; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.locator.PendingRangeMaps; +import org.openjdk.jmh.annotations.*; +import org.openjdk.jmh.infra.Blackhole; + +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.util.Collection; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.TimeUnit; + +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +@Warmup(iterations = 1, time = 1, timeUnit = TimeUnit.SECONDS) +@Measurement(iterations = 50, time = 1, timeUnit = TimeUnit.SECONDS) +@Fork(value = 3,jvmArgsAppend = "-Xmx512M") +@Threads(1) +@State(Scope.Benchmark) +public class PendingRangesBench +{ + PendingRangeMaps pendingRangeMaps; + int maxToken = 256 * 100; + + Multimap, InetAddress> oldPendingRanges; + + private Range genRange(String left, String right) + { + return new Range(new RandomPartitioner.BigIntegerToken(left), new RandomPartitioner.BigIntegerToken(right)); + } + + @Setup + public void setUp() throws UnknownHostException + { + pendingRangeMaps = new PendingRangeMaps(); + oldPendingRanges = HashMultimap.create(); + + InetAddress[] addresses = {InetAddress.getByName("127.0.0.1"), InetAddress.getByName("127.0.0.2")}; + + for (int i = 0; i < maxToken; i++) + { + for (int j = 0; j < ThreadLocalRandom.current().nextInt(2); j ++) + { + Range range = genRange(Integer.toString(i * 10 + 5), Integer.toString(i * 10 + 15)); + pendingRangeMaps.addPendingRange(range, addresses[j]); + oldPendingRanges.put(range, addresses[j]); + } + } + + // add the wrap around range + for (int j = 0; j < ThreadLocalRandom.current().nextInt(2); j ++) + { + Range range = genRange(Integer.toString(maxToken * 10 + 5), Integer.toString(5)); + pendingRangeMaps.addPendingRange(range, addresses[j]); + oldPendingRanges.put(range, addresses[j]); + } + } + + @Benchmark + public void searchToken(final Blackhole bh) + { + int randomToken = ThreadLocalRandom.current().nextInt(maxToken * 10 + 5); + Token searchToken = new RandomPartitioner.BigIntegerToken(Integer.toString(randomToken)); + bh.consume(pendingRangeMaps.pendingEndpointsFor(searchToken)); + } + + @Benchmark + public void searchTokenForOldPendingRanges(final Blackhole bh) + { + int randomToken = ThreadLocalRandom.current().nextInt(maxToken * 10 + 5); + Token searchToken = new RandomPartitioner.BigIntegerToken(Integer.toString(randomToken)); + Set endpoints = new HashSet<>(); + for (Map.Entry, Collection> entry : oldPendingRanges.asMap().entrySet()) + { + if (entry.getKey().contains(searchToken)) + endpoints.addAll(entry.getValue()); + } + bh.consume(endpoints); + } + +} diff --git a/test/unit/org/apache/cassandra/locator/PendingRangeMapsTest.java b/test/unit/org/apache/cassandra/locator/PendingRangeMapsTest.java new file mode 100644 index 0000000000..6d24447399 --- /dev/null +++ b/test/unit/org/apache/cassandra/locator/PendingRangeMapsTest.java @@ -0,0 +1,78 @@ +package org.apache.cassandra.locator; + +import org.apache.cassandra.dht.RandomPartitioner.BigIntegerToken; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.junit.Test; + +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.util.Collection; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class PendingRangeMapsTest { + + private Range genRange(String left, String right) + { + return new Range(new BigIntegerToken(left), new BigIntegerToken(right)); + } + + @Test + public void testPendingEndpoints() throws UnknownHostException + { + PendingRangeMaps pendingRangeMaps = new PendingRangeMaps(); + + pendingRangeMaps.addPendingRange(genRange("5", "15"), InetAddress.getByName("127.0.0.1")); + pendingRangeMaps.addPendingRange(genRange("15", "25"), InetAddress.getByName("127.0.0.2")); + pendingRangeMaps.addPendingRange(genRange("25", "35"), InetAddress.getByName("127.0.0.3")); + pendingRangeMaps.addPendingRange(genRange("35", "45"), InetAddress.getByName("127.0.0.4")); + pendingRangeMaps.addPendingRange(genRange("45", "55"), InetAddress.getByName("127.0.0.5")); + pendingRangeMaps.addPendingRange(genRange("45", "65"), InetAddress.getByName("127.0.0.6")); + + assertEquals(0, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("0")).size()); + assertEquals(0, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("5")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("10")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("15")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("20")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("25")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("35")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("45")).size()); + assertEquals(2, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("55")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("65")).size()); + + Collection endpoints = pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("15")); + assertTrue(endpoints.contains(InetAddress.getByName("127.0.0.1"))); + } + + @Test + public void testWrapAroundRanges() throws UnknownHostException + { + PendingRangeMaps pendingRangeMaps = new PendingRangeMaps(); + + pendingRangeMaps.addPendingRange(genRange("5", "15"), InetAddress.getByName("127.0.0.1")); + pendingRangeMaps.addPendingRange(genRange("15", "25"), InetAddress.getByName("127.0.0.2")); + pendingRangeMaps.addPendingRange(genRange("25", "35"), InetAddress.getByName("127.0.0.3")); + pendingRangeMaps.addPendingRange(genRange("35", "45"), InetAddress.getByName("127.0.0.4")); + pendingRangeMaps.addPendingRange(genRange("45", "55"), InetAddress.getByName("127.0.0.5")); + pendingRangeMaps.addPendingRange(genRange("45", "65"), InetAddress.getByName("127.0.0.6")); + pendingRangeMaps.addPendingRange(genRange("65", "7"), InetAddress.getByName("127.0.0.7")); + + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("0")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("5")).size()); + assertEquals(2, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("7")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("10")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("15")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("20")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("25")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("35")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("45")).size()); + assertEquals(2, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("55")).size()); + assertEquals(1, pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("65")).size()); + + Collection endpoints = pendingRangeMaps.pendingEndpointsFor(new BigIntegerToken("6")); + assertTrue(endpoints.contains(InetAddress.getByName("127.0.0.1"))); + assertTrue(endpoints.contains(InetAddress.getByName("127.0.0.7"))); + } +} From b551b8e1e6ac37698b78e4ee65a658bd446e7f05 Mon Sep 17 00:00:00 2001 From: Joel Knighton Date: Thu, 31 Dec 2015 10:27:25 -0600 Subject: [PATCH 06/16] Add check if existing fat client entry in gossip has same broadcast address in checkForEndpointCollision to enable quicker bootstrap retries. patch by jkni; reviewed by Stefania for CASSANDRA-10844 --- src/java/org/apache/cassandra/service/StorageService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 0698d11592..6e38b92b66 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -550,7 +550,7 @@ public class StorageService extends NotificationBroadcasterSupport implements IE for (Map.Entry entry : Gossiper.instance.getEndpointStates()) { - if (entry.getValue().getApplicationState(ApplicationState.STATUS) == null) + if (entry.getKey().equals(FBUtilities.getBroadcastAddress()) || entry.getValue().getApplicationState(ApplicationState.STATUS) == null) continue; String[] pieces = entry.getValue().getApplicationState(ApplicationState.STATUS).value.split(VersionedValue.DELIMITER_STR, -1); assert (pieces.length > 0); From c04bf2a9bb8bc48bcdf49455478836e6fd1f217d Mon Sep 17 00:00:00 2001 From: Ariel Weisberg Date: Tue, 29 Dec 2015 14:32:18 -0500 Subject: [PATCH 07/16] Enable GC logging by default patch by Chris Lohfink; reviewed by aweisberg for CASSANDRA-10140 --- CHANGES.txt | 1 + NEWS.txt | 2 ++ conf/cassandra-env.ps1 | 27 +++++++++--------- conf/cassandra-env.sh | 28 +++++++++---------- debian/patches/002cassandra_logdir_fix.dpatch | 18 ++++++++++-- 5 files changed, 45 insertions(+), 31 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 648200bbca..d5bb7a8453 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.2.5 + * Enable GC logging by default * Optimize pending range computation (CASSANDRA-9258) * Skip commit log and saved cache directories in SSTable version startup check (CASSANDRA-10902) * drop/alter user should be case sensitive (CASSANDRA-10817) diff --git a/NEWS.txt b/NEWS.txt index 3876c43060..57e321e83b 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -33,6 +33,8 @@ Operations "rack1". To override this behaviour use -Dcassandra.ignore_rack=true and/or -Dcassandra.ignore_dc=true. - Reloading the configuration file of GossipingPropertyFileSnitch has been disabled. + - GC logging is now enabled by default (but you can disable it if you want by + commenting the relevant lines of the cassandra-env file). New features ------------ diff --git a/conf/cassandra-env.ps1 b/conf/cassandra-env.ps1 index 970896407f..aff0d9e1b0 100644 --- a/conf/cassandra-env.ps1 +++ b/conf/cassandra-env.ps1 @@ -416,24 +416,23 @@ Function SetCassandraEnvironment $env:JVM_OPTS="$env:JVM_OPTS -XX:+CMSParallelInitialMarkEnabled -XX:+CMSEdenChunksRecordAlways" } - # GC logging options -- uncomment to enable - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintGCDetails" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintGCDateStamps" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintHeapAtGC" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintTenuringDistribution" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintGCApplicationStoppedTime" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintPromotionFailure" + # GC logging options + $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintGCDetails" + $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintGCDateStamps" + $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintHeapAtGC" + $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintTenuringDistribution" + $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintGCApplicationStoppedTime" + $env:JVM_OPTS="$env:JVM_OPTS -XX:+PrintPromotionFailure" # $env:JVM_OPTS="$env:JVM_OPTS -XX:PrintFLSStatistics=1" + + $env:JVM_OPTS="$env:JVM_OPTS -Xloggc:$env:CASSANDRA_HOME/logs/gc.log" + $env:JVM_OPTS="$env:JVM_OPTS -XX:+UseGCLogFileRotation" + $env:JVM_OPTS="$env:JVM_OPTS -XX:NumberOfGCLogFiles=10" + $env:JVM_OPTS="$env:JVM_OPTS -XX:GCLogFileSize=10M" + # if using version before JDK 6u34 or 7u2 use this instead of log rotation # $currentDate = (Get-Date).ToString('yyyy.MM.dd') # $env:JVM_OPTS="$env:JVM_OPTS -Xloggc:$env:CASSANDRA_HOME/logs/gc-$currentDate.log" - # If you are using JDK 6u34 7u2 or later you can enable GC log rotation - # don't stick the date in the log name if rotation is on. - # $env:JVM_OPTS="$env:JVM_OPTS -Xloggc:$env:CASSANDRA_HOME/logs/gc.log" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:+UseGCLogFileRotation" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:NumberOfGCLogFiles=10" - # $env:JVM_OPTS="$env:JVM_OPTS -XX:GCLogFileSize=10M" - # Configure the following for JEMallocAllocator and if jemalloc is not available in the system # library path. # set LD_LIBRARY_PATH=/lib/ diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh index e82198b3c0..ea1a73634c 100644 --- a/conf/cassandra-env.sh +++ b/conf/cassandra-env.sh @@ -232,21 +232,21 @@ if [ "$JVM_ARCH" = "64-Bit" ] ; then JVM_OPTS="$JVM_OPTS -XX:+UseCondCardMark" fi -# GC logging options -- uncomment to enable -# JVM_OPTS="$JVM_OPTS -XX:+PrintGCDetails" -# JVM_OPTS="$JVM_OPTS -XX:+PrintGCDateStamps" -# JVM_OPTS="$JVM_OPTS -XX:+PrintHeapAtGC" -# JVM_OPTS="$JVM_OPTS -XX:+PrintTenuringDistribution" -# JVM_OPTS="$JVM_OPTS -XX:+PrintGCApplicationStoppedTime" -# JVM_OPTS="$JVM_OPTS -XX:+PrintPromotionFailure" -# JVM_OPTS="$JVM_OPTS -XX:PrintFLSStatistics=1" +# GC logging options +JVM_OPTS="$JVM_OPTS -XX:+PrintGCDetails" +JVM_OPTS="$JVM_OPTS -XX:+PrintGCDateStamps" +JVM_OPTS="$JVM_OPTS -XX:+PrintHeapAtGC" +JVM_OPTS="$JVM_OPTS -XX:+PrintTenuringDistribution" +JVM_OPTS="$JVM_OPTS -XX:+PrintGCApplicationStoppedTime" +JVM_OPTS="$JVM_OPTS -XX:+PrintPromotionFailure" +#JVM_OPTS="$JVM_OPTS -XX:PrintFLSStatistics=1" + +JVM_OPTS="$JVM_OPTS -Xloggc:${CASSANDRA_HOME}/logs/gc.log" +JVM_OPTS="$JVM_OPTS -XX:+UseGCLogFileRotation" +JVM_OPTS="$JVM_OPTS -XX:NumberOfGCLogFiles=10" +JVM_OPTS="$JVM_OPTS -XX:GCLogFileSize=10M" +# if using version before JDK 6u34 or 7u2 use this instead of log rotation # JVM_OPTS="$JVM_OPTS -Xloggc:/var/log/cassandra/gc-`date +%s`.log" -# If you are using JDK 6u34 7u2 or later you can enable GC log rotation -# don't stick the date in the log name if rotation is on. -# JVM_OPTS="$JVM_OPTS -Xloggc:/var/log/cassandra/gc.log" -# JVM_OPTS="$JVM_OPTS -XX:+UseGCLogFileRotation" -# JVM_OPTS="$JVM_OPTS -XX:NumberOfGCLogFiles=10" -# JVM_OPTS="$JVM_OPTS -XX:GCLogFileSize=10M" # uncomment to have Cassandra JVM listen for remote debuggers/profilers on port 1414 # JVM_OPTS="$JVM_OPTS -agentlib:jdwp=transport=dt_socket,server=y,suspend=n,address=1414" diff --git a/debian/patches/002cassandra_logdir_fix.dpatch b/debian/patches/002cassandra_logdir_fix.dpatch index 8836eb46ff..cca337cca6 100644 --- a/debian/patches/002cassandra_logdir_fix.dpatch +++ b/debian/patches/002cassandra_logdir_fix.dpatch @@ -6,9 +6,9 @@ @DPATCH@ diff -urNad '--exclude=CVS' '--exclude=.svn' '--exclude=.git' '--exclude=.arch' '--exclude=.hg' '--exclude=_darcs' '--exclude=.bzr' cassandra~/bin/cassandra cassandra/bin/cassandra ---- cassandra~/bin/cassandra 2014-09-15 19:42:28.000000000 -0500 -+++ cassandra/bin/cassandra 2014-09-15 21:15:15.627505503 -0500 -@@ -134,7 +134,7 @@ +--- cassandra~/bin/cassandra 2015-10-27 14:15:10.718076265 -0500 ++++ cassandra/bin/cassandra 2015-10-27 14:23:10.000000000 -0500 +@@ -139,7 +139,7 @@ props="$3" class="$4" cassandra_parms="-Dlogback.configurationFile=logback.xml" @@ -17,3 +17,15 @@ diff -urNad '--exclude=CVS' '--exclude=.svn' '--exclude=.git' '--exclude=.arch' cassandra_parms="$cassandra_parms -Dcassandra.storagedir=$cassandra_storagedir" if [ "x$pidpath" != "x" ]; then +diff -urNad '--exclude=CVS' '--exclude=.svn' '--exclude=.git' '--exclude=.arch' '--exclude=.hg' '--exclude=_darcs' '--exclude=.bzr' cassandra~/conf/cassandra-env.sh cassandra/conf/cassandra-env.sh +--- cassandra~/conf/cassandra-env.sh 2015-10-27 14:20:22.990840135 -0500 ++++ cassandra/conf/cassandra-env.sh 2015-10-27 14:24:03.210202234 -0500 +@@ -288,7 +288,7 @@ + JVM_OPTS="$JVM_OPTS -XX:+PrintPromotionFailure" + #JVM_OPTS="$JVM_OPTS -XX:PrintFLSStatistics=1" + +-JVM_OPTS="$JVM_OPTS -Xloggc:${CASSANDRA_HOME}/logs/gc.log" ++JVM_OPTS="$JVM_OPTS -Xloggc:/var/log/cassandra/gc.log" + JVM_OPTS="$JVM_OPTS -XX:+UseGCLogFileRotation" + JVM_OPTS="$JVM_OPTS -XX:NumberOfGCLogFiles=10" + JVM_OPTS="$JVM_OPTS -XX:GCLogFileSize=10M" From 3e45fa1ab521bd50eb247f58daa2bfa76c6ab4e5 Mon Sep 17 00:00:00 2001 From: Ariel Weisberg Date: Tue, 29 Dec 2015 14:33:26 -0500 Subject: [PATCH 08/16] Enable GC logging by default (3.0 version) patch by Chris Lohfink; reviewed by aweisberg for CASSANDRA-10140 --- CHANGES.txt | 1 + NEWS.txt | 2 ++ conf/cassandra-env.ps1 | 3 +++ conf/cassandra-env.sh | 5 ++++- conf/jvm.options | 18 +++++++++--------- debian/patches/002cassandra_logdir_fix.dpatch | 18 +++++++++++++++--- 6 files changed, 34 insertions(+), 13 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 4b5610e6c5..103ae05ba5 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -11,6 +11,7 @@ * (Hadoop) Close Clusters and Sessions in Hadoop Input/Output classes (CASSANDRA-10837) * Fix sstableloader not working with upper case keyspace name (CASSANDRA-10806) Merged from 2.2: + * Enable GC logging by default (CASSANDRA-10140) * Optimize pending range computation (CASSANDRA-9258) * Skip commit log and saved cache directories in SSTable version startup check (CASSANDRA-10902) * drop/alter user should be case sensitive (CASSANDRA-10817) diff --git a/NEWS.txt b/NEWS.txt index 8a03e14717..26a83a90fc 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -32,6 +32,8 @@ Upgrading - Custom index implementation should be aware that the method Indexer::indexes() has been removed as its contract was misleading and all custom implementation should have almost surely returned true inconditionally for that method. + - GC logging is now enabled by default (you can disable it in the jvm.options + file if you prefer). 3.0 diff --git a/conf/cassandra-env.ps1 b/conf/cassandra-env.ps1 index a38429e1f5..5eefb04873 100644 --- a/conf/cassandra-env.ps1 +++ b/conf/cassandra-env.ps1 @@ -333,6 +333,9 @@ Function SetCassandraEnvironment ParseJVMInfo + #GC log path has to be defined here since it needs to find CASSANDRA_HOME + $env:JVM_OPTS="$env:JVM_OPTS -Xloggc:$env:CASSANDRA_HOME/logs/gc.log" + # Read user-defined JVM options from jvm.options file $content = Get-Content "$env:CASSANDRA_CONF\jvm.options" for ($i = 0; $i -lt $content.Count; $i++) diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh index ef164e8347..f32280384e 100644 --- a/conf/cassandra-env.sh +++ b/conf/cassandra-env.sh @@ -33,7 +33,7 @@ calculate_heap_sizes() Darwin) system_memory_in_bytes=`sysctl hw.memsize | awk '{print $2}'` system_memory_in_mb=`expr $system_memory_in_bytes / 1024 / 1024` - system_cpu_cores=`sysctl hw.ncpu | awk '{print $2}'` + ;; *) # assume reasonable defaults for e.g. a modern desktop or @@ -156,6 +156,9 @@ if [ "x$MALLOC_ARENA_MAX" = "x" ] ; then export MALLOC_ARENA_MAX=4 fi +#GC log path has to be defined here because it needs to access CASSANDRA_HOME +JVM_OPTS="$JVM_OPTS -Xloggc:${CASSANDRA_HOME}/logs/gc.log" + # Here we create the arguments that will get passed to the jvm when # starting cassandra. diff --git a/conf/jvm.options b/conf/jvm.options index c5d3d95ed1..a7b3bd87f3 100644 --- a/conf/jvm.options +++ b/conf/jvm.options @@ -95,14 +95,14 @@ ### GC logging options -- uncomment to enable -#-XX:+PrintGCDetails -#-XX:+PrintGCDateStamps -#-XX:+PrintHeapAtGC -#-XX:+PrintTenuringDistribution -#-XX:+PrintGCApplicationStoppedTime -#-XX:+PrintPromotionFailure +-XX:+PrintGCDetails +-XX:+PrintGCDateStamps +-XX:+PrintHeapAtGC +-XX:+PrintTenuringDistribution +-XX:+PrintGCApplicationStoppedTime +-XX:+PrintPromotionFailure #-XX:PrintFLSStatistics=1 #-Xloggc:/var/log/cassandra/gc.log -#-XX:+UseGCLogFileRotation -#-XX:NumberOfGCLogFiles=10 -#-XX:GCLogFileSize=10M +-XX:+UseGCLogFileRotation +-XX:NumberOfGCLogFiles=10 +-XX:GCLogFileSize=10M diff --git a/debian/patches/002cassandra_logdir_fix.dpatch b/debian/patches/002cassandra_logdir_fix.dpatch index 8836eb46ff..87387b9c67 100644 --- a/debian/patches/002cassandra_logdir_fix.dpatch +++ b/debian/patches/002cassandra_logdir_fix.dpatch @@ -6,9 +6,9 @@ @DPATCH@ diff -urNad '--exclude=CVS' '--exclude=.svn' '--exclude=.git' '--exclude=.arch' '--exclude=.hg' '--exclude=_darcs' '--exclude=.bzr' cassandra~/bin/cassandra cassandra/bin/cassandra ---- cassandra~/bin/cassandra 2014-09-15 19:42:28.000000000 -0500 -+++ cassandra/bin/cassandra 2014-09-15 21:15:15.627505503 -0500 -@@ -134,7 +134,7 @@ +--- cassandra~/bin/cassandra 2015-10-27 14:35:22.000000000 -0500 ++++ cassandra/bin/cassandra 2015-10-27 14:41:38.000000000 -0500 +@@ -139,7 +139,7 @@ props="$3" class="$4" cassandra_parms="-Dlogback.configurationFile=logback.xml" @@ -17,3 +17,15 @@ diff -urNad '--exclude=CVS' '--exclude=.svn' '--exclude=.git' '--exclude=.arch' cassandra_parms="$cassandra_parms -Dcassandra.storagedir=$cassandra_storagedir" if [ "x$pidpath" != "x" ]; then +diff -urNad '--exclude=CVS' '--exclude=.svn' '--exclude=.git' '--exclude=.arch' '--exclude=.hg' '--exclude=_darcs' '--exclude=.bzr' cassandra~/conf/cassandra-env.sh cassandra/conf/cassandra-env.sh +--- cassandra~/conf/cassandra-env.sh 2015-10-27 14:40:39.000000000 -0500 ++++ cassandra/conf/cassandra-env.sh 2015-10-27 14:42:40.647449856 -0500 +@@ -204,7 +204,7 @@ + esac + + #GC log path has to be defined here because it needs to access CASSANDRA_HOME +-JVM_OPTS="$JVM_OPTS -Xloggc:${CASSANDRA_HOME}/logs/gc.log" ++JVM_OPTS="$JVM_OPTS -Xloggc:/var/log/cassandra/gc.log" + + # Here we create the arguments that will get passed to the jvm when + # starting cassandra. From 9ca7e162667dbabddfc2fba10eeb82bae5bdf56e Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Tue, 5 Jan 2016 18:04:21 +0100 Subject: [PATCH 09/16] Minor CHANGES.txt fixup --- CHANGES.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index d5bb7a8453..fc87c7d636 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,5 @@ 2.2.5 - * Enable GC logging by default + * Enable GC logging by default (CASSANDRA-10140) * Optimize pending range computation (CASSANDRA-9258) * Skip commit log and saved cache directories in SSTable version startup check (CASSANDRA-10902) * drop/alter user should be case sensitive (CASSANDRA-10817) From 12b6c0a5264138a6f569e5e957362ddaf1a4a0be Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Wed, 6 Jan 2016 10:44:22 +0100 Subject: [PATCH 10/16] Fix typo in javadoc --- src/java/org/apache/cassandra/db/rows/Row.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/java/org/apache/cassandra/db/rows/Row.java b/src/java/org/apache/cassandra/db/rows/Row.java index 8a67e9b3fe..5f79a66552 100644 --- a/src/java/org/apache/cassandra/db/rows/Row.java +++ b/src/java/org/apache/cassandra/db/rows/Row.java @@ -394,11 +394,11 @@ public interface Row extends Unfiltered, Collection public Clustering clustering(); /** - * Adds the liveness information for the partition key columns of this row. + * Adds the liveness information for the primary key columns of this row. * * This call is optional (skipping it is equivalent to calling {@code addPartitionKeyLivenessInfo(LivenessInfo.NONE)}). * - * @param info the liveness information for the partition key columns of the built row. + * @param info the liveness information for the primary key columns of the built row. */ public void addPrimaryKeyLivenessInfo(LivenessInfo info); From f937c8bee87666ac1ba1ac96e0c7803629cc3508 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Wed, 6 Jan 2016 14:44:29 +0100 Subject: [PATCH 11/16] Revert wrong line removal from bad merge --- conf/cassandra-env.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh index f32280384e..83fe4c5443 100644 --- a/conf/cassandra-env.sh +++ b/conf/cassandra-env.sh @@ -33,7 +33,7 @@ calculate_heap_sizes() Darwin) system_memory_in_bytes=`sysctl hw.memsize | awk '{print $2}'` system_memory_in_mb=`expr $system_memory_in_bytes / 1024 / 1024` - + system_cpu_cores=`sysctl hw.ncpu | awk '{print $2}'` ;; *) # assume reasonable defaults for e.g. a modern desktop or From 70c08ece563731cd546d24541f225c888f4d02f5 Mon Sep 17 00:00:00 2001 From: Carl Yeksigian Date: Wed, 6 Jan 2016 10:41:47 -0500 Subject: [PATCH 12/16] MV timestamp should be the maximum of the values, not the minimum patch by Carl Yeksigian; reviewed by Jake Luciani for CASSANDRA-10910 --- .../apache/cassandra/db/view/TemporalRow.java | 65 +++++++++++++------ .../org/apache/cassandra/cql3/ViewTest.java | 30 +++++++++ 2 files changed, 76 insertions(+), 19 deletions(-) diff --git a/src/java/org/apache/cassandra/db/view/TemporalRow.java b/src/java/org/apache/cassandra/db/view/TemporalRow.java index 8898857183..8ee310da6b 100644 --- a/src/java/org/apache/cassandra/db/view/TemporalRow.java +++ b/src/java/org/apache/cassandra/db/view/TemporalRow.java @@ -22,6 +22,7 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.Comparator; import java.util.HashMap; import java.util.Iterator; import java.util.List; @@ -279,9 +280,7 @@ public class TemporalRow this.nowInSec = nowInSec; LivenessInfo liveness = row.primaryKeyLivenessInfo(); - this.viewClusteringLocalDeletionTime = minValueIfSet(viewClusteringLocalDeletionTime, row.deletion().time().localDeletionTime(), NO_DELETION_TIME); - this.viewClusteringTimestamp = minValueIfSet(viewClusteringTimestamp, liveness.timestamp(), NO_TIMESTAMP); - this.viewClusteringTtl = minValueIfSet(viewClusteringTtl, liveness.ttl(), NO_TTL); + updateLiveness(liveness.ttl(), liveness.timestamp(), row.deletion().time().localDeletionTime()); List clusteringDefs = baseCfs.metadata.clusteringColumns(); clusteringColumns = new HashMap<>(); @@ -295,6 +294,31 @@ public class TemporalRow } } + /* + * PK ts:5, ttl:1, deletion: 2 + * Col ts:4, ttl:2, deletion: 3 + * + * TTL use min, since it expires at the lowest time which we are expiring. If we have the above values, we + * would want to return 1, since the base row expires in 1 second. + * + * Timestamp uses max, as this is the time that the row has been written to the view. See CASSANDRA-10910. + * + * Local Deletion Time should use max, as this deletion will cover all previous values written. + */ + @SuppressWarnings("unchecked") + private void updateLiveness(int ttl, long timestamp, int localDeletionTime) + { + // We are returning whichever is higher from valueIfSet + // Natural order will return the max: 1.compareTo(2) < 0, so 2 is returned + // Reverse order will return the min: 1.compareTo(2) > 0, so 1 is returned + final Comparator max = Comparator.naturalOrder(); + final Comparator min = Comparator.reverseOrder(); + + this.viewClusteringTtl = valueIfSet(viewClusteringTtl, ttl, NO_TTL, min); + this.viewClusteringTimestamp = valueIfSet(viewClusteringTimestamp, timestamp, NO_TIMESTAMP, max); + this.viewClusteringLocalDeletionTime = valueIfSet(viewClusteringLocalDeletionTime, localDeletionTime, NO_DELETION_TIME, max); + } + @Override public String toString() { @@ -351,30 +375,33 @@ public class TemporalRow // If this column is part of the view's primary keys if (viewPrimaryKey.contains(identifier)) { - this.viewClusteringTtl = minValueIfSet(this.viewClusteringTtl, ttl, NO_TTL); - this.viewClusteringTimestamp = minValueIfSet(this.viewClusteringTimestamp, timestamp, NO_TIMESTAMP); - this.viewClusteringLocalDeletionTime = minValueIfSet(this.viewClusteringLocalDeletionTime, localDeletionTime, NO_DELETION_TIME); + updateLiveness(ttl, timestamp, localDeletionTime); } innerMap.get(cellPath).setVersion(new TemporalCell(value, timestamp, ttl, localDeletionTime, isNew)); } - private static int minValueIfSet(int existing, int update, int defaultValue) + /** + * @return + *
    + *
  • + * If both existing and update are defaultValue, return defaultValue + *
  • + *
  • + * If only one of existing or existing are defaultValue, return the one which is not + *
  • + *
  • + * If both existing and update are not defaultValue, compare using comparator and return the higher one. + *
  • + *
+ */ + private static T valueIfSet(T existing, T update, T defaultValue, Comparator comparator) { - if (existing == defaultValue) + if (existing.equals(defaultValue)) return update; - if (update == defaultValue) + if (update.equals(defaultValue)) return existing; - return Math.min(existing, update); - } - - private static long minValueIfSet(long existing, long update, long defaultValue) - { - if (existing == defaultValue) - return update; - if (update == defaultValue) - return existing; - return Math.min(existing, update); + return comparator.compare(existing, update) > 0 ? existing : update; } public int viewClusteringTtl() diff --git a/test/unit/org/apache/cassandra/cql3/ViewTest.java b/test/unit/org/apache/cassandra/cql3/ViewTest.java index 8ae21df199..2e3cf5f2cb 100644 --- a/test/unit/org/apache/cassandra/cql3/ViewTest.java +++ b/test/unit/org/apache/cassandra/cql3/ViewTest.java @@ -274,6 +274,36 @@ public class ViewTest extends CQLTester assertRows(execute("SELECT c from mv_tstest where k = 0 and val = ?", "baz"), row(1)); } + @Test + public void testRegularColumnTimestampUpdates() throws Throwable + { + // Regression test for CASSANDRA-10910 + + createTable("CREATE TABLE %s (" + + "k int PRIMARY KEY, " + + "c int, " + + "val int)"); + + execute("USE " + keyspace()); + executeNet(protocolVersion, "USE " + keyspace()); + + createView("mv_rctstest", "CREATE MATERIALIZED VIEW %s AS SELECT * FROM %%s WHERE k IS NOT NULL AND c IS NOT NULL PRIMARY KEY (k,c)"); + + updateView("UPDATE %s SET c = ?, val = ? WHERE k = ?", 0, 0, 0); + updateView("UPDATE %s SET val = ? WHERE k = ?", 1, 0); + updateView("UPDATE %s SET c = ? WHERE k = ?", 1, 0); + assertRows(execute("SELECT c, k, val FROM mv_rctstest"), row(1, 0, 1)); + + updateView("TRUNCATE %s"); + + updateView("UPDATE %s USING TIMESTAMP 1 SET c = ?, val = ? WHERE k = ?", 0, 0, 0); + updateView("UPDATE %s USING TIMESTAMP 3 SET c = ? WHERE k = ?", 1, 0); + updateView("UPDATE %s USING TIMESTAMP 2 SET val = ? WHERE k = ?", 1, 0); + updateView("UPDATE %s USING TIMESTAMP 4 SET c = ? WHERE k = ?", 2, 0); + updateView("UPDATE %s USING TIMESTAMP 3 SET val = ? WHERE k = ?", 2, 0); + assertRows(execute("SELECT c, k, val FROM mv_rctstest"), row(2, 0, 2)); + } + @Test public void testCountersTable() throws Throwable { From e9e127abf7915d1351dac5a4c008c35ecfebf9ef Mon Sep 17 00:00:00 2001 From: Carl Yeksigian Date: Wed, 6 Jan 2016 11:53:10 -0500 Subject: [PATCH 13/16] changes.txt --- CHANGES.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGES.txt b/CHANGES.txt index 103ae05ba5..cf872d9f5f 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 3.0.3 + * MV should use the maximum timestamp of the primary key (CASSANDRA-10910) * Fix potential assertion error during compaction (CASSANDRA-10944) * Fix counting of received sstables in streaming (CASSANDRA-10949) * Implement hints compression (CASSANDRA-9428) From 11716547f87f4c88a2790323744b0ed97175854d Mon Sep 17 00:00:00 2001 From: Stefania Alborghetti Date: Fri, 20 Nov 2015 13:51:06 +0800 Subject: [PATCH 14/16] Match cassandra-loader options in COPY FROM patch by Stefania; reviewed by pauloricardomg for CASSANDRA-9303 --- CHANGES.txt | 1 + NEWS.txt | 7 + bin/cqlsh | 133 +- conf/cqlshrc.sample | 17 +- pylib/cqlshlib/copyutil.py | 1253 +++++++++++++---- pylib/cqlshlib/formatting.py | 96 +- .../cql3/statements/BatchStatement.java | 27 +- .../cassandra/transport/ServerConnection.java | 2 +- tools/bin/cassandra-stress.bat | 2 +- 9 files changed, 1128 insertions(+), 410 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 0bd5485ceb..844a28f873 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.1.13 + * Match cassandra-loader options in COPY FROM (CASSANDRA-9303) * Fix binding to any address in CqlBulkRecordWriter (CASSANDRA-9309) * Fix the way we replace sstables after anticompaction (CASSANDRA-10831) * cqlsh fails to decode utf-8 characters for text typed columns (CASSANDRA-10875) diff --git a/NEWS.txt b/NEWS.txt index 7a15d0245e..088efae766 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -13,6 +13,13 @@ restore snapshots created with the previous major version using the 'sstableloader' tool. You can upgrade the file format of your snapshots using the provided 'sstableupgrade' tool. +2.1.13 +====== + +New features +------------ + - New options for cqlsh COPY FROM and COPY TO, see CASSANDRA-9303 for details. + 2.1.12 ====== diff --git a/bin/cqlsh b/bin/cqlsh index 1d5db257c8..e2bfcd7379 100755 --- a/bin/cqlsh +++ b/bin/cqlsh @@ -41,7 +41,6 @@ import optparse import os import platform import sys -import time import traceback import warnings @@ -119,7 +118,8 @@ cqlshlibdir = os.path.join(CASSANDRA_PATH, 'pylib') if os.path.isdir(cqlshlibdir): sys.path.insert(0, cqlshlibdir) -from cqlshlib import cql3handling, cqlhandling, copyutil, pylexotron, sslhandling +from cqlshlib import cql3handling, cqlhandling, pylexotron, sslhandling +from cqlshlib.copyutil import ExportTask, ImportTask from cqlshlib.displaying import (ANSI_RESET, BLUE, COLUMN_NAME_COLORS, CYAN, RED, FormattedValue, colorme) from cqlshlib.formatting import (format_by_type, format_value_utype, @@ -410,10 +410,12 @@ def complete_copy_column_names(ctxt, cqlsh): return set(colnames[1:]) - set(existcols) -COPY_COMMON_OPTIONS = ['DELIMITER', 'QUOTE', 'ESCAPE', 'HEADER', 'NULL', - 'MAXATTEMPTS', 'REPORTFREQUENCY'] -COPY_FROM_OPTIONS = ['CHUNKSIZE', 'INGESTRATE', 'MAXBATCHSIZE', 'MINBATCHSIZE'] -COPY_TO_OPTIONS = ['ENCODING', 'TIMEFORMAT', 'PAGESIZE', 'PAGETIMEOUT', 'MAXREQUESTS'] +COPY_COMMON_OPTIONS = ['DELIMITER', 'QUOTE', 'ESCAPE', 'HEADER', 'NULL', 'DATETIMEFORMAT', + 'MAXATTEMPTS', 'REPORTFREQUENCY', 'DECIMALSEP', 'THOUSANDSSEP', 'BOOLSTYLE', + 'NUMPROCESSES', 'CONFIGFILE', 'RATEFILE'] +COPY_FROM_OPTIONS = ['CHUNKSIZE', 'INGESTRATE', 'MAXBATCHSIZE', 'MINBATCHSIZE', 'MAXROWS', + 'SKIPROWS', 'SKIPCOLS', 'MAXPARSEERRORS', 'MAXINSERTERRORS', 'ERRFILE'] +COPY_TO_OPTIONS = ['ENCODING', 'PAGESIZE', 'PAGETIMEOUT', 'BEGINTOKEN', 'ENDTOKEN', 'MAXOUTPUTSIZE', 'MAXREQUESTS'] @cqlsh_syntax_completer('copyOption', 'optnames') @@ -521,23 +523,6 @@ warnings.showwarning = show_warning_without_quoting_line warnings.filterwarnings('always', category=cql3handling.UnexpectedTableStructure) -def describe_interval(seconds): - desc = [] - for length, unit in ((86400, 'day'), (3600, 'hour'), (60, 'minute')): - num = int(seconds) / length - if num > 0: - desc.append('%d %s' % (num, unit)) - if num > 1: - desc[-1] += 's' - seconds %= length - words = '%.03f seconds' % seconds - if len(desc) > 1: - words = ', '.join(desc) + ', and ' + words - elif len(desc) == 1: - words = desc[0] + ' and ' + words - return words - - def insert_driver_hooks(): extend_cql_deserialization() auto_format_udts() @@ -604,8 +589,7 @@ class Shell(cmd.Cmd): last_hist = None shunted_query_out = None use_paging = True - csv_dialect_defaults = dict(delimiter=',', doublequote=False, - escapechar='\\', quotechar='"') + default_page_size = 100 def __init__(self, hostname, port, color=False, @@ -1509,32 +1493,67 @@ class Shell(cmd.Cmd): COPY x TO: Exports data from a Cassandra table in CSV format. COPY [ ( column [, ...] ) ] - FROM ( '' | STDIN ) + FROM ( '' | STDIN ) [ WITH