From cb1f1399b139029e5b1c12a4bf65d19a55724933 Mon Sep 17 00:00:00 2001 From: Cameron Zemek Date: Wed, 13 Sep 2023 14:41:50 +1000 Subject: [PATCH] Improve performance of compactions when table does not have an index patch by Cameron Zemek; reviewed by Branimir Lambov, Stefan Miklosovic for CASSANDRA-18773 --- CHANGES.txt | 1 + .../db/compaction/CompactionIterator.java | 13 +++- .../UnfilteredPartitionIterators.java | 65 ++++++++++++++----- 3 files changed, 63 insertions(+), 16 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 1ede25e409..88d734026a 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0.12 + * Improve performance of compactions when table does not have an index (CASSANDRA-18773) * JMH improvements - faster build and async profiler (CASSANDRA-18871) * Enable 3rd party JDK installations for Debian package (CASSANDRA-18844) * Fix NTS log message when an unrecognized strategy option is passed (CASSANDRA-18679) diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index ec6a4d464c..39b28a85c4 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -154,6 +154,17 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte { return new UnfilteredPartitionIterators.MergeListener() { + private boolean rowProcessingNeeded() + { + return type == OperationType.COMPACTION && controller.cfs.indexManager.hasIndexes(); + } + + @Override + public boolean preserveOrder() + { + return rowProcessingNeeded(); + } + public UnfilteredRowIterators.MergeListener getRowMergeListener(DecoratedKey partitionKey, List versions) { int merged = 0; @@ -169,7 +180,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte CompactionIterator.this.updateCounterFor(merged); - if (type != OperationType.COMPACTION || !controller.cfs.indexManager.hasIndexes()) + if (!rowProcessingNeeded()) return null; Columns statics = Columns.NONE; diff --git a/src/java/org/apache/cassandra/db/partitions/UnfilteredPartitionIterators.java b/src/java/org/apache/cassandra/db/partitions/UnfilteredPartitionIterators.java index a051ee1a9d..fe3e74b72b 100644 --- a/src/java/org/apache/cassandra/db/partitions/UnfilteredPartitionIterators.java +++ b/src/java/org/apache/cassandra/db/partitions/UnfilteredPartitionIterators.java @@ -45,10 +45,30 @@ public abstract class UnfilteredPartitionIterators public interface MergeListener { + /** + * Returns true if the merger needs to preserve the position of sources within the merge when passing data to + * the listener. If false, the merger can avoid creating empty sources for non-present partitions and + * significantly speed up processing. + * + * @return True to preserve position of source iterators. + */ + public default boolean preserveOrder() { return true; } public UnfilteredRowIterators.MergeListener getRowMergeListener(DecoratedKey partitionKey, List versions); public default void close() {} - public static MergeListener NOOP = (partitionKey, versions) -> UnfilteredRowIterators.MergeListener.NOOP; + public static MergeListener NOOP = new MergeListener() + { + @Override + public boolean preserveOrder() + { + return false; + } + + public UnfilteredRowIterators.MergeListener getRowMergeListener(DecoratedKey partitionKey, List versions) + { + return UnfilteredRowIterators.MergeListener.NOOP; + } + }; } @SuppressWarnings("resource") // The created resources are returned right away @@ -108,6 +128,8 @@ public abstract class UnfilteredPartitionIterators final TableMetadata metadata = iterators.get(0).metadata(); + final boolean preserveOrder = listener != null && listener.preserveOrder(); + final MergeIterator merged = MergeIterator.get(iterators, partitionComparator, new MergeIterator.Reducer() { private final List toMerge = new ArrayList<>(iterators.size()); @@ -120,9 +142,16 @@ public abstract class UnfilteredPartitionIterators partitionKey = current.partitionKey(); isReverseOrder = current.isReverseOrder(); - // Note that because the MergeListener cares about it, we want to preserve the index of the iterator. - // Non-present iterator will thus be set to empty in getReduced. - toMerge.set(idx, current); + if (preserveOrder) + { + // Note that because the MergeListener cares about it, we want to preserve the index of the iterator. + // Non-present iterator will thus be set to empty in getReduced. + toMerge.set(idx, current); + } + else + { + toMerge.add(current); + } } @SuppressWarnings("resource") @@ -132,17 +161,20 @@ public abstract class UnfilteredPartitionIterators ? null : listener.getRowMergeListener(partitionKey, toMerge); - // Make a single empty iterator object to merge, we don't need toMerge.size() copiess - UnfilteredRowIterator empty = null; - - // Replace nulls by empty iterators - for (int i = 0; i < toMerge.size(); i++) + if (preserveOrder) { - if (toMerge.get(i) == null) + // Make a single empty iterator object to merge, we don't need toMerge.size() copiess + UnfilteredRowIterator empty = null; + + // Replace nulls by empty iterators + for (int i = 0; i < toMerge.size(); i++) { - if (null == empty) - empty = EmptyIterators.unfilteredRow(metadata, partitionKey, isReverseOrder); - toMerge.set(i, empty); + if (toMerge.get(i) == null) + { + if (null == empty) + empty = EmptyIterators.unfilteredRow(metadata, partitionKey, isReverseOrder); + toMerge.set(i, empty); + } } } @@ -152,8 +184,11 @@ public abstract class UnfilteredPartitionIterators protected void onKeyChange() { toMerge.clear(); - for (int i = 0; i < iterators.size(); i++) - toMerge.add(null); + if (preserveOrder) + { + for (int i = 0; i < iterators.size(); i++) + toMerge.add(null); + } } });