diff --git a/CHANGES.txt b/CHANGES.txt index c368787bd6..ac6c53a488 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -18,6 +18,7 @@ * SSTable metadata(Stats.db) format change (CASSANDRA-6356) * Push composites support in the storage engine (CASSANDRA-5417) * Add snapshot space used to cfstats (CASSANDRA-6231) + * Add cardinality estimator for key count estimation (CASSANDRA-5906) 2.0.4 diff --git a/lib/licenses/stream-2.5.1.txt b/lib/licenses/stream-2.5.1.txt new file mode 100644 index 0000000000..c8dc6770e5 --- /dev/null +++ b/lib/licenses/stream-2.5.1.txt @@ -0,0 +1,202 @@ + + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + APPENDIX: How to apply the Apache License to your work. + + To apply the Apache License to your work, attach the following + boilerplate notice, with the fields enclosed by brackets "[]" + replaced with your own identifying information. (Don't include + the brackets!) The text should be enclosed in the appropriate + comment syntax for the file format. We also recommend that a + file or class name and description of purpose be included on the + same "printed page" as the copyright notice for easier + identification within third-party archives. + + Copyright 2011 Clearspring Technologies + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. diff --git a/lib/stream-2.5.1.jar b/lib/stream-2.5.1.jar new file mode 100644 index 0000000000..17f0014827 Binary files /dev/null and b/lib/stream-2.5.1.jar differ diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index b72f91cb76..e4f5237f8a 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -550,7 +550,7 @@ public class CompactionManager implements CompactionManagerMBean long totalkeysWritten = 0; int expectedBloomFilterSize = Math.max(cfs.metadata.getIndexInterval(), - (int) (SSTableReader.getApproximateKeyCount(Arrays.asList(sstable), cfs.metadata))); + (int) (SSTableReader.getApproximateKeyCount(Arrays.asList(sstable)))); if (logger.isDebugEnabled()) logger.debug("Expected bloom filter size : {}", expectedBloomFilterSize); diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java index cabe4865cd..61f98f0fdd 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java @@ -118,7 +118,7 @@ public class CompactionTask extends AbstractCompactionTask long start = System.nanoTime(); long totalkeysWritten = 0; - long estimatedTotalKeys = Math.max(cfs.metadata.getIndexInterval(), SSTableReader.getApproximateKeyCount(actuallyCompact, cfs.metadata)); + long estimatedTotalKeys = Math.max(cfs.metadata.getIndexInterval(), SSTableReader.getApproximateKeyCount(actuallyCompact)); long estimatedSSTables = Math.max(1, SSTableReader.getTotalBytes(actuallyCompact) / strategy.getMaxSSTableBytes()); long keysPerSSTable = (long) Math.ceil((double) estimatedTotalKeys / estimatedSSTables); if (logger.isDebugEnabled()) diff --git a/src/java/org/apache/cassandra/db/compaction/Scrubber.java b/src/java/org/apache/cassandra/db/compaction/Scrubber.java index bec29d5a9b..eabfdbc7a6 100644 --- a/src/java/org/apache/cassandra/db/compaction/Scrubber.java +++ b/src/java/org/apache/cassandra/db/compaction/Scrubber.java @@ -85,7 +85,7 @@ public class Scrubber implements Closeable ? new ScrubController(cfs) : new CompactionController(cfs, Collections.singleton(sstable), CompactionManager.getDefaultGcBefore(cfs)); this.isCommutative = cfs.metadata.getDefaultValidator().isCommutative(); - this.expectedBloomFilterSize = Math.max(cfs.metadata.getIndexInterval(), (int)(SSTableReader.getApproximateKeyCount(toScrub,cfs.metadata))); + this.expectedBloomFilterSize = Math.max(cfs.metadata.getIndexInterval(), (int)(SSTableReader.getApproximateKeyCount(toScrub))); // loop through each row, deserializing to check for damage. // we'll also loop through the index at the same time, using the position from the index to recover if the diff --git a/src/java/org/apache/cassandra/db/compaction/Upgrader.java b/src/java/org/apache/cassandra/db/compaction/Upgrader.java index e4d29e9c56..de966685de 100644 --- a/src/java/org/apache/cassandra/db/compaction/Upgrader.java +++ b/src/java/org/apache/cassandra/db/compaction/Upgrader.java @@ -56,7 +56,7 @@ public class Upgrader this.controller = new UpgradeController(cfs); this.strategy = cfs.getCompactionStrategy(); - long estimatedTotalKeys = Math.max(cfs.metadata.getIndexInterval(), SSTableReader.getApproximateKeyCount(toUpgrade, cfs.metadata)); + long estimatedTotalKeys = Math.max(cfs.metadata.getIndexInterval(), SSTableReader.getApproximateKeyCount(toUpgrade)); long estimatedSSTables = Math.max(1, SSTableReader.getTotalBytes(this.toUpgrade) / strategy.getMaxSSTableBytes()); this.estimatedRows = (long) Math.ceil((double) estimatedTotalKeys / estimatedSSTables); } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java index de877bccdf..30bfd77cb9 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java @@ -20,16 +20,16 @@ package org.apache.cassandra.io.sstable; import java.io.*; import java.nio.ByteBuffer; import java.util.*; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.ScheduledThreadPoolExecutor; -import java.util.concurrent.TimeUnit; +import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; +import com.clearspring.analytics.stream.cardinality.CardinalityMergeException; +import com.clearspring.analytics.stream.cardinality.ICardinality; import com.google.common.annotations.VisibleForTesting; +import com.google.common.base.Predicate; +import com.google.common.collect.Iterators; import com.google.common.collect.Ordering; import com.google.common.primitives.Longs; import com.google.common.util.concurrent.RateLimiter; @@ -46,7 +46,6 @@ import org.apache.cassandra.db.commitlog.ReplayPosition; import org.apache.cassandra.db.compaction.ICompactionScanner; import org.apache.cassandra.db.index.SecondaryIndex; import org.apache.cassandra.dht.*; -import org.apache.cassandra.io.FSReadError; import org.apache.cassandra.io.compress.CompressedRandomAccessReader; import org.apache.cassandra.io.compress.CompressedThrottledReader; import org.apache.cassandra.io.compress.CompressionMetadata; @@ -134,18 +133,69 @@ public class SSTableReader extends SSTable implements Closeable public RestorableMeter readMeter; private ScheduledFuture readMeterSyncFuture; - public static long getApproximateKeyCount(Iterable sstables, CFMetaData metadata) + /** + * Calculate approximate key count. + * If cardinality estimator is available on all given sstables, then this method use them to estimate + * key count. + * If not, then this uses index summaries. + * + * @param sstables SSTables to calculate key count + * @return estimated key count + */ + public static long getApproximateKeyCount(Collection sstables) { - long count = 0; + long count = -1; - for (SSTableReader sstable : sstables) + // check if cardinality estimator is available for all SSTables + boolean cardinalityAvailable = !sstables.isEmpty() && Iterators.all(sstables.iterator(), new Predicate() { - // using getMaxIndexSummarySize() lets us ignore the current sampling level - count += (sstable.getMaxIndexSummarySize() + 1) * sstable.indexSummary.getSamplingLevel(); - if (logger.isDebugEnabled()) - logger.debug("index size for bloom filter calc for file : {} : {}", sstable.getFilename(), count); + public boolean apply(SSTableReader sstable) + { + return sstable.descriptor.version.newStatsFile; + } + }); + + // if it is, load them to estimate key count + if (cardinalityAvailable) + { + boolean failed = false; + ICardinality cardinality = null; + for (SSTableReader sstable : sstables) + { + try + { + CompactionMetadata metadata = (CompactionMetadata) sstable.descriptor.getMetadataSerializer().deserialize(sstable.descriptor, MetadataType.COMPACTION); + if (cardinality == null) + cardinality = metadata.cardinalityEstimator; + else + cardinality = cardinality.merge(metadata.cardinalityEstimator); + } + catch (IOException e) + { + logger.warn("Reading cardinality from Statistics.db failed.", e); + failed = true; + break; + } + catch (CardinalityMergeException e) + { + logger.warn("Cardinality merge failed.", e); + failed = true; + break; + } + } + if (cardinality != null && !failed) + count = cardinality.cardinality(); } + // if something went wrong above or cardinality is not available, calculate using index summary + if (count < 0) + { + for (SSTableReader sstable : sstables) + { + // using getMaxIndexSummarySize() lets us ignore the current sampling level + count += (sstable.getMaxIndexSummarySize() + 1) * sstable.indexSummary.getSamplingLevel(); + } + } return count; } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java index ac8e2b2807..7b034283a2 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java @@ -149,6 +149,7 @@ public class SSTableWriter extends SSTable private void afterAppend(DecoratedKey decoratedKey, long dataPosition, RowIndexEntry index) { + sstableMetadataCollector.addKey(decoratedKey.key); lastWrittenKey = decoratedKey; last = lastWrittenKey; if (first == null) diff --git a/src/java/org/apache/cassandra/io/sstable/metadata/CompactionMetadata.java b/src/java/org/apache/cassandra/io/sstable/metadata/CompactionMetadata.java index fd0e62615a..1dd33e8cf2 100644 --- a/src/java/org/apache/cassandra/io/sstable/metadata/CompactionMetadata.java +++ b/src/java/org/apache/cassandra/io/sstable/metadata/CompactionMetadata.java @@ -23,8 +23,12 @@ import java.io.IOException; import java.util.HashSet; import java.util.Set; +import com.clearspring.analytics.stream.cardinality.HyperLogLogPlus; +import com.clearspring.analytics.stream.cardinality.ICardinality; + import org.apache.cassandra.db.TypeSizes; import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.utils.ByteBufferUtil; /** * Compaction related SSTable metadata. @@ -37,9 +41,12 @@ public class CompactionMetadata extends MetadataComponent public final Set ancestors; - public CompactionMetadata(Set ancestors) + public final ICardinality cardinalityEstimator; + + public CompactionMetadata(Set ancestors, ICardinality cardinalityEstimator) { this.ancestors = ancestors; + this.cardinalityEstimator = cardinalityEstimator; } public MetadataType getType() @@ -71,6 +78,8 @@ public class CompactionMetadata extends MetadataComponent size += TypeSizes.NATIVE.sizeof(component.ancestors.size()); for (int g : component.ancestors) size += TypeSizes.NATIVE.sizeof(g); + byte[] serializedCardinality = component.cardinalityEstimator.getBytes(); + size += TypeSizes.NATIVE.sizeof(serializedCardinality.length) + serializedCardinality.length; return size; } @@ -79,6 +88,7 @@ public class CompactionMetadata extends MetadataComponent out.writeInt(component.ancestors.size()); for (int g : component.ancestors) out.writeInt(g); + ByteBufferUtil.writeWithLength(component.cardinalityEstimator.getBytes(), out); } public CompactionMetadata deserialize(Descriptor.Version version, DataInput in) throws IOException @@ -87,7 +97,8 @@ public class CompactionMetadata extends MetadataComponent Set ancestors = new HashSet<>(nbAncestors); for (int i = 0; i < nbAncestors; i++) ancestors.add(in.readInt()); - return new CompactionMetadata(ancestors); + ICardinality cardinality = HyperLogLogPlus.Builder.build(ByteBufferUtil.readBytes(in, in.readInt())); + return new CompactionMetadata(ancestors, cardinality); } } } diff --git a/src/java/org/apache/cassandra/io/sstable/metadata/LegacyMetadataSerializer.java b/src/java/org/apache/cassandra/io/sstable/metadata/LegacyMetadataSerializer.java index a691591ac6..33d4f16966 100644 --- a/src/java/org/apache/cassandra/io/sstable/metadata/LegacyMetadataSerializer.java +++ b/src/java/org/apache/cassandra/io/sstable/metadata/LegacyMetadataSerializer.java @@ -148,7 +148,7 @@ public class LegacyMetadataSerializer extends MetadataSerializer maxColumnNames)); if (types.contains(MetadataType.COMPACTION)) components.put(MetadataType.COMPACTION, - new CompactionMetadata(ancestors)); + new CompactionMetadata(ancestors, null)); } } return components; diff --git a/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java b/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java index c125a98c7a..e20015d160 100644 --- a/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java +++ b/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java @@ -21,12 +21,15 @@ import java.io.File; import java.nio.ByteBuffer; import java.util.*; +import com.clearspring.analytics.stream.cardinality.HyperLogLogPlus; +import com.clearspring.analytics.stream.cardinality.ICardinality; import com.google.common.collect.Maps; import org.apache.cassandra.db.commitlog.ReplayPosition; import org.apache.cassandra.db.composites.CellNameType; import org.apache.cassandra.io.sstable.*; import org.apache.cassandra.utils.EstimatedHistogram; +import org.apache.cassandra.utils.MurmurHash; import org.apache.cassandra.utils.StreamingHistogram; public class MetadataCollector @@ -77,6 +80,13 @@ public class MetadataCollector protected int sstableLevel; protected List minColumnNames = Collections.emptyList(); protected List maxColumnNames = Collections.emptyList(); + /** + * Default cardinality estimation method is to use HyperLogLog++. + * Parameter here(p=13, sp=25) should give reasonable estimation + * while lowering bytes required to hold information. + * See CASSANDRA-5906 for detail. + */ + protected ICardinality cardinality = new HyperLogLogPlus(13, 25); private final CellNameType columnNameComparator; public MetadataCollector(CellNameType columnNameComparator) @@ -103,6 +113,12 @@ public class MetadataCollector } } + public void addKey(ByteBuffer key) + { + long hashed = MurmurHash.hash2_64(key, key.position(), key.remaining(), 0); + cardinality.offerHashed(hashed); + } + public void addRowSize(long rowSize) { estimatedRowSize.add(rowSize); @@ -213,7 +229,7 @@ public class MetadataCollector sstableLevel, minColumnNames, maxColumnNames)); - components.put(MetadataType.COMPACTION, new CompactionMetadata(ancestors)); + components.put(MetadataType.COMPACTION, new CompactionMetadata(ancestors, cardinality)); return components; } diff --git a/src/java/org/apache/cassandra/tools/SSTableMetadataViewer.java b/src/java/org/apache/cassandra/tools/SSTableMetadataViewer.java index d8166ad07b..a2f7b89ad5 100644 --- a/src/java/org/apache/cassandra/tools/SSTableMetadataViewer.java +++ b/src/java/org/apache/cassandra/tools/SSTableMetadataViewer.java @@ -66,6 +66,11 @@ public class SSTableMetadataViewer out.println(stats.replayPosition); printHistograms(stats, out); } + if (compaction != null) + { + out.printf("Ancestors: %s%n", compaction.ancestors.toString()); + out.printf("Estimated cardinality: %s%n", compaction.cardinalityEstimator.cardinality()); + } } }