diff --git a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java index c8340a6b3c..6847437f2a 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java +++ b/src/java/org/apache/cassandra/service/accord/AccordCommandStore.java @@ -116,7 +116,6 @@ public class AccordCommandStore extends CommandStore this.lock = lock; } - @Override public AccordSafeCommand acquireIfLoaded(TxnId txnId) { diff --git a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java index 6797b59ec0..9710cc34d9 100644 --- a/src/java/org/apache/cassandra/service/accord/AccordExecutor.java +++ b/src/java/org/apache/cassandra/service/accord/AccordExecutor.java @@ -543,6 +543,11 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo waitingToRun(task, task.executor()); } + void submitExclusive(Runnable run) + { + submitExclusive(new PlainRunnable(run)); + } + void submitExclusive(AccordTask task) { ++tasks; @@ -1246,6 +1251,11 @@ public abstract class AccordExecutor implements CacheSize, AccordCacheEntry.OnLo final Runnable run; final @Nullable SequentialExecutor executor; + PlainRunnable(Runnable run) + { + this(null, run, null); + } + PlainRunnable(AsyncPromise result, Runnable run, @Nullable SequentialExecutor executor) { this.result = result; diff --git a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java index ff43778d69..d51b297932 100644 --- a/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java +++ b/src/java/org/apache/cassandra/service/accord/CommandsForRanges.java @@ -18,12 +18,18 @@ package org.apache.cassandra.service.accord; +import java.util.ArrayList; import java.util.Comparator; +import java.util.HashMap; import java.util.Iterator; +import java.util.List; import java.util.Map; import java.util.NavigableMap; import java.util.TreeMap; import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.atomic.AtomicReferenceFieldUpdater; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.function.BiFunction; import java.util.function.Consumer; import java.util.function.UnaryOperator; @@ -54,7 +60,9 @@ import accord.utils.SymmetricComparator; import accord.utils.UnhandledEnum; import org.agrona.collections.Object2ObjectHashMap; import org.apache.cassandra.service.accord.api.TokenKey; +import org.apache.cassandra.utils.btree.BTreeSet; import org.apache.cassandra.utils.btree.IntervalBTree; +import org.apache.cassandra.utils.concurrent.IntrusiveStack; import static accord.local.CommandSummaries.SummaryStatus.NOT_DIRECTLY_WITNESSED; import static org.apache.cassandra.utils.btree.IntervalBTree.InclusiveEndHelper.endWithStart; @@ -122,8 +130,43 @@ public class CommandsForRanges extends TreeMap implements Co return this; } - public static class Manager implements AccordCache.Listener + public static class Manager implements AccordCache.Listener, Runnable { + static class IntervalTreeEdit extends IntrusiveStack + { + final TxnId txnId; + final @Nullable Object[] update, remove; + + IntervalTreeEdit(TxnId txnId, Object[] update, Object[] remove) + { + this.txnId = txnId; + this.update = update; + this.remove = remove; + } + + public static boolean push(IntervalTreeEdit edit, Manager manager) + { + return null == IntrusiveStack.getAndPush(pendingEditsUpdater, manager, edit); + } + + public IntervalTreeEdit reverse() + { + return reverse(this); + } + + boolean isSize(int size) + { + return IntrusiveStack.isSize(size, this); + } + + IntervalTreeEdit merge(IntervalTreeEdit next) + { + Invariants.require(this.txnId.equals(next.txnId)); + Object[] remove = this.remove == null ? next.remove : next.remove == null ? this.remove : IntervalBTree.update(this.remove, next.remove, COMPARATORS); + return new IntervalTreeEdit(txnId, next.update, remove); + } + } + private final AccordCommandStore commandStore; private final RangeSearcher searcher; private final AtomicReference> transitive = new AtomicReference<>(new TreeMap<>()); @@ -131,6 +174,10 @@ public class CommandsForRanges extends TreeMap implements Co private final Object2ObjectHashMap cachedRangeTxnsById = new Object2ObjectHashMap<>(); private Object[] cachedRangeTxnsByRange = IntervalBTree.empty(); + private volatile IntervalTreeEdit pendingEdits; + private final Lock drainPendingEditsLock = new ReentrantLock(); + private static final AtomicReferenceFieldUpdater pendingEditsUpdater = AtomicReferenceFieldUpdater.newUpdater(Manager.class, IntervalTreeEdit.class, "pendingEdits"); + public Manager(AccordCommandStore commandStore) { this.commandStore = commandStore; @@ -155,16 +202,93 @@ public class CommandsForRanges extends TreeMap implements Co { RangeRoute cur = cachedRangeTxnsById.put(cmd.txnId(), upd); if (!upd.equals(cur)) - { - if (cur != null) - remove(txnId, cur); - cachedRangeTxnsByRange = IntervalBTree.update(cachedRangeTxnsByRange, toMap(txnId, upd), COMPARATORS); - } + pushEdit(new IntervalTreeEdit(txnId, toMap(txnId, upd), cur == null ? null : toMap(txnId, cur))); } } } } + private void pushEdit(IntervalTreeEdit edit) + { + if (IntervalTreeEdit.push(edit, this)) + commandStore.executor().submitExclusive(this); + } + + @Override + public void run() + { + if (drainPendingEditsLock.tryLock()) + { + try + { + drainPendingEditsInternal(); + } + finally + { + drainPendingEditsLock.unlock(); + postUnlock(); + } + } + } + + Object[] cachedRangeTxnsByRange() + { + drainPendingEditsLock.lock(); + try + { + drainPendingEditsInternal(); + return cachedRangeTxnsByRange; + } + finally + { + drainPendingEditsLock.unlock(); + postUnlock(); + } + } + + void drainPendingEditsInternal() + { + IntervalTreeEdit edits = pendingEditsUpdater.getAndSet(this, null); + if (edits == null) + return; + + if (edits.isSize(1)) + { + if (edits.remove != null) cachedRangeTxnsByRange = IntervalBTree.subtract(cachedRangeTxnsByRange, edits.remove, COMPARATORS); + if (edits.update != null) cachedRangeTxnsByRange = IntervalBTree.update(cachedRangeTxnsByRange, edits.update, COMPARATORS); + return; + } + + edits = edits.reverse(); + Map editMap = new HashMap<>(); + for (IntervalTreeEdit edit : edits) + editMap.merge(edit.txnId, edit, IntervalTreeEdit::merge); + + List update = new ArrayList<>(), remove = new ArrayList<>(); + for (IntervalTreeEdit edit : editMap.values()) + { + if (edit.update != null) update.addAll(BTreeSet.wrap(edit.update, COMPARATORS.totalOrder())); + if (edit.remove != null) remove.addAll(BTreeSet.wrap(edit.remove, COMPARATORS.totalOrder())); + } + + if (!remove.isEmpty()) + { + remove.sort(COMPARATORS.totalOrder()); + cachedRangeTxnsByRange = IntervalBTree.subtract(cachedRangeTxnsByRange, IntervalBTree.build(remove, COMPARATORS), COMPARATORS); + } + if (!update.isEmpty()) + { + update.sort(COMPARATORS.totalOrder()); + cachedRangeTxnsByRange = IntervalBTree.update(cachedRangeTxnsByRange, IntervalBTree.build(update, COMPARATORS), COMPARATORS); + } + } + + private void postUnlock() + { + if (pendingEdits != null) + commandStore.executor().submit(this); + } + @Override public void onEvict(AccordCacheEntry state) { @@ -173,15 +297,10 @@ public class CommandsForRanges extends TreeMap implements Co { RangeRoute cur = cachedRangeTxnsById.remove(txnId); if (cur != null) - remove(txnId, cur); + pushEdit(new IntervalTreeEdit(txnId, null, toMap(txnId, cur))); } } - private void remove(TxnId txnId, RangeRoute route) - { - cachedRangeTxnsByRange = IntervalBTree.subtract(cachedRangeTxnsByRange, toMap(txnId, route), COMPARATORS); - } - static Object[] toMap(TxnId txnId, RangeRoute route) { int size = route.size(); @@ -199,7 +318,6 @@ public class CommandsForRanges extends TreeMap implements Co } } } - } public CommandsForRanges.Loader loader(@Nullable TxnId primaryTxnId, KeyHistory keyHistory, Unseekables keysOrRanges) @@ -331,7 +449,7 @@ public class CommandsForRanges extends TreeMap implements Co { for (RoutingKey key : (AbstractUnseekableKeys)keysOrRanges) { - IntervalBTree.accumulate(manager.cachedRangeTxnsByRange, KEY_COMPARATORS, key, (f, s, i, c) -> { + IntervalBTree.accumulate(manager.cachedRangeTxnsByRange(), KEY_COMPARATORS, key, (f, s, i, c) -> { TxnIdInterval interval = (TxnIdInterval)i; if (isRelevant(interval)) { @@ -349,7 +467,7 @@ public class CommandsForRanges extends TreeMap implements Co { for (Range range : (AbstractRanges)keysOrRanges) { - IntervalBTree.accumulate(manager.cachedRangeTxnsByRange, COMPARATORS, new TxnIdInterval(range.start(), range.end(), TxnId.NONE), (f, s, i, c) -> { + IntervalBTree.accumulate(manager.cachedRangeTxnsByRange(), COMPARATORS, new TxnIdInterval(range.start(), range.end(), TxnId.NONE), (f, s, i, c) -> { if (isRelevant(i)) { TxnId txnId = i.txnId; diff --git a/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java b/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java index 7ea8b81989..b027295562 100644 --- a/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java +++ b/src/java/org/apache/cassandra/utils/btree/IntervalBTree.java @@ -20,6 +20,7 @@ package org.apache.cassandra.utils.btree; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collection; import java.util.Comparator; import java.util.List; @@ -595,6 +596,16 @@ public class IntervalBTree } } + public static Object[] build(Collection build, IntervalComparators comparators) + { + try (FastIntervalTreeBuilder builder = IntervalBTree.fastBuilder(comparators)) + { + for (V v : build) + builder.add(v); + return builder.build(); + } + } + /** * Build a tree of unknown size, in order. */ diff --git a/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java b/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java index d27ed598f0..0dcc3b509a 100644 --- a/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java +++ b/src/java/org/apache/cassandra/utils/concurrent/IntrusiveStack.java @@ -163,6 +163,13 @@ public class IntrusiveStack> implements Iterable return size; } + protected static boolean isSize(int size, IntrusiveStack list) + { + while (list != null && --size >= 0) + list = list.next; + return list == null && size == 0; + } + protected static > long accumulate(T list, LongAccumulator accumulator, long initialValue) { long value = initialValue;