diff --git a/mantis-network/build.gradle b/mantis-network/build.gradle index 314091388..06e8df7dd 100644 --- a/mantis-network/build.gradle +++ b/mantis-network/build.gradle @@ -14,6 +14,20 @@ * limitations under the License. */ +buildscript { + repositories { + gradlePluginPortal() + } + dependencies { + classpath 'me.champeau.jmh:jmh-gradle-plugin:0.7.3' + } +} + +// JMH microbenchmarks live in src/jmh. Run with: ./gradlew :mantis-network:jmh +// or build the self-contained jar (./gradlew :mantis-network:jmhJar) and drive it with the +// standard JMH CLI, which is the only way to vary -t (thread count). +apply plugin: 'me.champeau.jmh' + ext { mqlVersion = '3.4.+' nettyVersion = '4.1.17.Final' @@ -30,8 +44,19 @@ dependencies { testImplementation libraries.junitJupiter testImplementation libraries.mockitoCore testImplementation libraries.slf4jLog4j12 + + // The benchmarks drive the routers and the push-server batch wrap directly, so they need the + // metrics implementation that main only compiles against. + jmhImplementation libraries.spectatorApi } test { useJUnitPlatform() } + +jmh { + jmhVersion = '1.37' + fork = 1 + warmupIterations = 5 + iterations = 5 +} diff --git a/mantis-network/src/jmh/java/io/reactivex/mantis/network/push/LegacyRoundRobinRouter.java b/mantis-network/src/jmh/java/io/reactivex/mantis/network/push/LegacyRoundRobinRouter.java new file mode 100644 index 000000000..c170503fc --- /dev/null +++ b/mantis-network/src/jmh/java/io/reactivex/mantis/network/push/LegacyRoundRobinRouter.java @@ -0,0 +1,106 @@ +/* + * Copyright 2019 Netflix, Inc. + * + * 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. + */ + +package io.reactivex.mantis.network.push; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.Iterator; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.Map.Entry; +import java.util.Set; +import java.util.concurrent.atomic.AtomicReference; +import rx.functions.Func1; + +/** + * The pre-change {@link RoundRobinRouter}, kept as the JMH baseline arm. + * + *

{@code route} and {@code loopingIterator} below are the verbatim bodies from before the + * change, copied out of git rather than reimplemented, precisely so the comparison cannot drift: + * + *

+ *   git show master:mantis-network/src/main/java/io/reactivex/mantis/network/push/RoundRobinRouter.java
+ * 
+ * + *

Do not "clean this up" or apply review suggestions to it — its only value is being byte-for-byte + * what shipped. + */ +public class LegacyRoundRobinRouter extends Router { + + public LegacyRoundRobinRouter(String name, Func1 encoder) { + super("LegacyRoundRobinRouter_" + name, encoder); + } + + @Override + public void route(Set> connections, List chunks) { + if (chunks != null && !chunks.isEmpty()) { + numEventsProcessed.increment(chunks.size()); + } + List> randomOrder = new ArrayList<>(connections); + Collections.shuffle(randomOrder); + if (chunks != null && !chunks.isEmpty() && !randomOrder.isEmpty()) { + Iterator> iter = loopingIterator(randomOrder); + Map, List> writes = new HashMap<>(); + // process chunks + for (T chunk : chunks) { + AsyncConnection connection = iter.next(); + Func1 predicate = connection.getPredicate(); + if (predicate == null || predicate.call(chunk)) { + List buffer = writes.get(connection); + if (buffer == null) { + buffer = new LinkedList<>(); + writes.put(connection, buffer); + } + buffer.add(encoder.call(chunk)); + } + } + if (!writes.isEmpty()) { + for (Entry, List> entry : writes.entrySet()) { + AsyncConnection connection = entry.getKey(); + List toWrite = entry.getValue(); + connection.write(toWrite); + numEventsRouted.increment(toWrite.size()); + } + } + } + } + + private Iterator> loopingIterator(final Collection> connections) { + final AtomicReference>> iterRef = new AtomicReference<>(connections.iterator()); + return + new Iterator>() { + @Override + public boolean hasNext() { + return true; + } + + @Override + public AsyncConnection next() { + Iterator> iter = iterRef.get(); + if (iter.hasNext()) { + return iter.next(); + } else { + iterRef.set(connections.iterator()); + return iterRef.get().next(); + } + } + }; + } +} diff --git a/mantis-network/src/jmh/java/io/reactivex/mantis/network/push/RoundRobinRouterBenchmark.java b/mantis-network/src/jmh/java/io/reactivex/mantis/network/push/RoundRobinRouterBenchmark.java new file mode 100644 index 000000000..30a388b26 --- /dev/null +++ b/mantis-network/src/jmh/java/io/reactivex/mantis/network/push/RoundRobinRouterBenchmark.java @@ -0,0 +1,167 @@ +/* + * Copyright 2019 Netflix, Inc. + * + * 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. + */ + +package io.reactivex.mantis.network.push; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.Warmup; +import rx.Observer; +import rx.functions.Func1; + +/** + * Compares {@link RoundRobinRouter#route} against {@link LegacyRoundRobinRouter#route}, which is the + * verbatim pre-change body. + * + *

The two arms differ in three places: the legacy one copies the connection {@link Set} into an + * {@link ArrayList} and {@link java.util.Collections#shuffle(List) shuffle}s it, buffers per-destination + * writes in a {@link java.util.LinkedList}, and allocates an unsized {@link java.util.HashMap}. The new + * one advances a looping iterator to a single {@code ThreadLocalRandom} offset and sizes both + * collections. + * + *

Run multi-threaded. The largest claimed cost is not the per-connection work, it is that + * {@code Collections.shuffle(List)} draws from a single {@code private static Random} shared by the + * whole JVM, so every draw is a contended CAS on one seed. That only shows up with more than one + * router thread, and JMH's thread count is not expressible as a {@code @Param}, so drive it from the + * CLI: + * + *

+ *   ./gradlew :mantis-network:jmhJar
+ *   java -jar mantis-network/build/libs/mantis-network-*-jmh.jar RoundRobinRouterBenchmark -t 1
+ *   java -jar mantis-network/build/libs/mantis-network-*-jmh.jar RoundRobinRouterBenchmark -t 8
+ * 
+ */ +@State(Scope.Benchmark) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.MICROSECONDS) +@Warmup(iterations = 5, time = 1) +@Measurement(iterations = 5, time = 1) +@Fork(1) +public class RoundRobinRouterBenchmark { + + private static final AtomicInteger NAMES = new AtomicInteger(); + + /** Subscribers on one group. Real groups run from a handful to hundreds of thousands. */ + @Param({"2", "8", "64", "512"}) + public int connections; + + /** Events in one drain -- what the 200ms chunker hands to route() at a time. */ + @Param({"30", "200"}) + public int chunkSize; + + private Router current; + private Router legacy; + private Set> connectionSet; + private List chunks; + + /** + * Counts deliveries into a plain field on a per-connection object. Not volatile and not + * synchronized on purpose: each connection has its own sink, so there is nothing to contend on, + * and the store still cannot be optimised away because the sink stays reachable from the + * benchmark state. Deliberately not a Counter -- {@code CounterImpl.value()} reads through to a + * Spectator registry that is a no-op outside a running worker, so it always reads zero here. + */ + static final class Sink implements Observer> { + long batches; + long events; + + @Override + public void onNext(List data) { + batches++; + events += data.size(); + } + + @Override + public void onCompleted() { + } + + @Override + public void onError(Throwable e) { + } + } + + private List sinks; + + @Setup(Level.Trial) + public void setup() { + String tag = "bench" + NAMES.incrementAndGet(); + Func1 encoder = b -> b; + current = new RoundRobinRouter<>(tag, encoder); + legacy = new LegacyRoundRobinRouter<>(tag, encoder); + + sinks = new ArrayList<>(connections); + connectionSet = new HashSet<>(connections * 2); + for (int i = 0; i < connections; i++) { + Sink sink = new Sink(); + sinks.add(sink); + connectionSet.add(new AsyncConnection<>( + "host" + i, 7000 + i, "id" + i, "slot" + i, "group", sink, null)); + } + + chunks = new ArrayList<>(chunkSize); + for (int i = 0; i < chunkSize; i++) { + byte[] event = new byte[256]; + event[0] = (byte) i; + chunks.add(event); + } + + // Cross-check the arms deliver the same volume before measuring anything. Content and + // destination differ run to run -- both arms randomise the starting connection -- so the + // invariant that holds is every chunk goes out exactly once. + long before = totalEvents(); + current.route(connectionSet, chunks); + long afterCurrent = totalEvents() - before; + legacy.route(connectionSet, chunks); + long afterLegacy = totalEvents() - before - afterCurrent; + if (afterCurrent != chunkSize || afterLegacy != chunkSize) { + throw new IllegalStateException("arms disagree: current=" + afterCurrent + + " legacy=" + afterLegacy + " expected=" + chunkSize); + } + } + + private long totalEvents() { + long total = 0; + for (Sink sink : sinks) { + total += sink.events; + } + return total; + } + + @Benchmark + public void route() { + current.route(connectionSet, chunks); + } + + @Benchmark + public void routeLegacyVerbatim() { + legacy.route(connectionSet, chunks); + } +} diff --git a/mantis-network/src/main/java/io/reactivex/mantis/network/push/RoundRobinRouter.java b/mantis-network/src/main/java/io/reactivex/mantis/network/push/RoundRobinRouter.java index e2cbf3083..438c4f2c9 100644 --- a/mantis-network/src/main/java/io/reactivex/mantis/network/push/RoundRobinRouter.java +++ b/mantis-network/src/main/java/io/reactivex/mantis/network/push/RoundRobinRouter.java @@ -18,14 +18,13 @@ import java.util.ArrayList; import java.util.Collection; -import java.util.Collections; import java.util.HashMap; import java.util.Iterator; -import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.Set; +import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.atomic.AtomicReference; import rx.functions.Func1; @@ -38,36 +37,51 @@ public RoundRobinRouter(String name, Func1 encoder) { @Override public void route(Set> connections, List chunks) { - if (chunks != null && !chunks.isEmpty()) { - numEventsProcessed.increment(chunks.size()); + if (chunks == null || chunks.isEmpty()) { + return; } - List> randomOrder = new ArrayList<>(connections); - Collections.shuffle(randomOrder); - if (chunks != null && !chunks.isEmpty() && !randomOrder.isEmpty()) { - Iterator> iter = loopingIterator(randomOrder); - Map, List> writes = new HashMap<>(); - // process chunks - for (T chunk : chunks) { - AsyncConnection connection = iter.next(); - Func1 predicate = connection.getPredicate(); - if (predicate == null || predicate.call(chunk)) { - List buffer = writes.get(connection); - if (buffer == null) { - buffer = new LinkedList<>(); - writes.put(connection, buffer); - } - buffer.add(encoder.call(chunk)); - } - } - if (!writes.isEmpty()) { - for (Entry, List> entry : writes.entrySet()) { - AsyncConnection connection = entry.getKey(); - List toWrite = entry.getValue(); - connection.write(toWrite); - numEventsRouted.increment(toWrite.size()); + numEventsProcessed.increment(chunks.size()); + + int numConnections = connections.size(); + if (numConnections == 0) { + return; + } + + // Start the round robin at a random offset rather than shuffling the whole connection set. + // A shuffle costs one Random.nextInt() and one swap per connection, and Collections.shuffle + // (List) draws from a single private static Random shared by the whole JVM, so every draw is + // a contended CAS on that one seed -- shared with every other caller of shuffle() in the + // process. A single ThreadLocalRandom draw gives the same uniform expectation per + // connection with no shared state and no per-connection work, which matters here because + // route() is called once per subscriber group per drain. + Iterator> iter = loopingIterator(connections); + for (int i = ThreadLocalRandom.current().nextInt(numConnections); i > 0; i--) { + iter.next(); + } + + // assume even distribution + int bufferCapacity = (chunks.size() / numConnections) + 1; + Map, List> writes = + new HashMap<>(Math.min(numConnections, chunks.size())); + // process chunks + for (T chunk : chunks) { + AsyncConnection connection = iter.next(); + Func1 predicate = connection.getPredicate(); + if (predicate == null || predicate.call(chunk)) { + List buffer = writes.get(connection); + if (buffer == null) { + buffer = new ArrayList<>(bufferCapacity); + writes.put(connection, buffer); } + buffer.add(encoder.call(chunk)); } } + for (Entry, List> entry : writes.entrySet()) { + AsyncConnection connection = entry.getKey(); + List toWrite = entry.getValue(); + connection.write(toWrite); + numEventsRouted.increment(toWrite.size()); + } } private Iterator> loopingIterator(final Collection> connections) { diff --git a/mantis-network/src/test/java/io/reactivex/mantis/network/push/RoundRobinRouterTest.java b/mantis-network/src/test/java/io/reactivex/mantis/network/push/RoundRobinRouterTest.java new file mode 100644 index 000000000..713ad3d20 --- /dev/null +++ b/mantis-network/src/test/java/io/reactivex/mantis/network/push/RoundRobinRouterTest.java @@ -0,0 +1,200 @@ +/* + * Copyright 2026 Netflix, Inc. + * + * 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. + */ + +package io.reactivex.mantis.network.push; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.LinkedHashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import org.junit.jupiter.api.Test; +import rx.Observer; +import rx.functions.Func1; + +/** + * Behavioural pins for {@link RoundRobinRouter#route}. The change under test replaced a full + * {@code Collections.shuffle} of the connection set with a single random start offset, and swapped + * the per-destination {@code LinkedList} for a pre-sized {@code ArrayList}; neither is observable + * from the outside, so what these tests assert is the routing contract that must not move. + */ +class RoundRobinRouterTest { + + /** Records what a connection was handed, in order. */ + private static final class Recorder implements Observer> { + + private final List> writes = new ArrayList<>(); + + @Override + public void onNext(List data) { + writes.add(data.stream() + .map(b -> new String(b, StandardCharsets.UTF_8)) + .collect(Collectors.toList())); + } + + @Override + public void onError(Throwable e) { + throw new AssertionError(e); + } + + @Override + public void onCompleted() { + } + + List received() { + return writes.stream().flatMap(List::stream).collect(Collectors.toList()); + } + } + + private static RoundRobinRouter router(String name) { + return new RoundRobinRouter<>(name, s -> s.getBytes(StandardCharsets.UTF_8)); + } + + /** Connections keyed by id, insertion-ordered so failures are readable. */ + private static Map connect( + Set> into, int count, Func1 predicate) { + Map recorders = new LinkedHashMap<>(); + for (int i = 0; i < count; i++) { + String id = "conn-" + i; + Recorder recorder = new Recorder(); + recorders.put(id, recorder); + into.add(new AsyncConnection<>("host", 1000 + i, id, id, "group", recorder, predicate)); + } + return recorders; + } + + private static List chunks(int count) { + return IntStream.range(0, count).mapToObj(i -> "event-" + i).collect(Collectors.toList()); + } + + @Test + void everyChunkIsDeliveredExactlyOnceAndSpreadEvenly() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 4, null); + List chunks = chunks(10); + + RoundRobinRouter router = router("even-spread"); + router.route(connections, chunks); + + List delivered = recorders.values().stream() + .flatMap(r -> r.received().stream()) + .collect(Collectors.toList()); + assertEquals(chunks.size(), delivered.size(), "no chunk may be dropped or duplicated"); + assertEquals(new HashSet<>(chunks), new HashSet<>(delivered)); + + // 10 chunks over 4 connections is 3 or 2 each, whatever the start offset. + for (Map.Entry e : recorders.entrySet()) { + int size = e.getValue().received().size(); + assertTrue(size == 2 || size == 3, e.getKey() + " received " + size); + } + } + + /** + * Each destination gets one write holding all of its chunks -- not one write per chunk. This is + * what makes the downstream batching worthwhile, so it is worth pinning. + */ + @Test + void eachDestinationReceivesASingleBatchedWrite() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 2, null); + + router("batched").route(connections, chunks(6)); + + for (Map.Entry e : recorders.entrySet()) { + assertEquals(1, e.getValue().writes.size(), e.getKey() + " should get one write"); + assertEquals(3, e.getValue().received().size()); + } + } + + /** Fewer chunks than connections: only the chunks.size() connections in line get a write. */ + @Test + void connectionsWithNoChunksAreNotWrittenTo() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 10, null); + + router("sparse").route(connections, chunks(3)); + + long written = recorders.values().stream().filter(r -> !r.writes.isEmpty()).count(); + assertEquals(3, written, "exactly one write per chunk, to distinct connections"); + } + + /** + * A chunk assigned to a connection whose predicate rejects it is dropped, not handed on to the + * next connection. That was the pre-existing behaviour and this change does not alter it. + */ + @Test + void chunksRejectedByAPredicateAreDroppedNotReRouted() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 3, s -> s.endsWith("0")); + List chunks = chunks(9); + + RoundRobinRouter router = router("predicate"); + router.route(connections, chunks); + + List delivered = recorders.values().stream() + .flatMap(r -> r.received().stream()) + .collect(Collectors.toList()); + assertEquals(List.of("event-0"), delivered); + } + + /** + * The start offset must keep moving: with a single chunk and many connections, load only spreads + * if successive calls begin somewhere else. Over 2000 calls across 50 connections the chance of + * any connection being missed by a uniform draw is about 1e-16, so this is not flaky. + */ + @Test + void startOffsetIsRandomisedAcrossCallsSoASingleChunkSpreads() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 50, null); + + RoundRobinRouter router = router("spread"); + for (int i = 0; i < 2000; i++) { + router.route(connections, List.of("event-" + i)); + } + + long touched = recorders.values().stream().filter(r -> !r.writes.isEmpty()).count(); + assertEquals(50, touched, "every connection should have been the start at least once"); + } + + @Test + void emptyAndNullChunkBatchesAreNoOps() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 3, null); + + RoundRobinRouter router = router("empty"); + router.route(connections, List.of()); + router.route(connections, null); + + assertTrue(recorders.values().stream().allMatch(r -> r.writes.isEmpty())); + } + + @Test + void noConnectionsDropsTheBatchWithoutThrowing() { + Set> connections = new LinkedHashSet<>(); + Map recorders = connect(connections, 2, null); + // deliberately route to an empty set, not to `connections` + router("no-connections").route(new HashSet<>(), chunks(5)); + assertTrue(recorders.values().stream().allMatch(r -> r.writes.isEmpty())); + } +}