Skip to content

Commit 430fd92

Browse files
committed
Removed benchVertx from Multi-Producer Suites (Mpmc, Mpsc)
1 parent f98eaf3 commit 430fd92

4 files changed

Lines changed: 13 additions & 56 deletions

File tree

‎zthread-benchmark/src/main/java/io/github/namanoncode/zthread/benchmark/adapters/VertxAdapter.java‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@ public void start(EventHandler handler, int consumers) {
1919
this.vertx = Vertx.vertx(options);
2020
this.eventBus = vertx.eventBus();
2121

22-
eventBus.localConsumer(ADDRESS, message -> {
22+
io.vertx.core.eventbus.MessageConsumer<Object> consumer = eventBus.localConsumer(ADDRESS);
23+
consumer.setMaxBufferedMessages(10_000_000);
24+
consumer.handler(message -> {
2325
handler.onEvent((BenchmarkEvent) message.body());
2426
});
2527
}

‎zthread-benchmark/src/main/java/io/github/namanoncode/zthread/benchmark/throughput/MpmcEventBenchmark.java‎

Lines changed: 6 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -110,18 +110,13 @@ public void teardownInvocation() {
110110
activeConcurrentQueue = null;
111111
}
112112

113-
private void runProducers(Runnable producerTask) throws InterruptedException {
114-
CountDownLatch producersDone = new CountDownLatch(PRODUCERS);
115-
for (int p = 0; p < PRODUCERS; p++) {
116-
producerPool.submit(() -> {
117-
try {
118-
producerTask.run();
119-
} finally {
120-
producersDone.countDown();
121-
}
122-
});
113+
protected void runProducers(Runnable producerTask) throws InterruptedException {
114+
java.util.concurrent.ExecutorService executor = java.util.concurrent.Executors.newFixedThreadPool(PRODUCERS);
115+
for (int i = 0; i < PRODUCERS; i++) {
116+
executor.submit(producerTask);
123117
}
124-
producersDone.await();
118+
executor.shutdown();
119+
executor.awaitTermination(1, java.util.concurrent.TimeUnit.MINUTES);
125120
}
126121

127122
private void startBlockingQueueConsumers(BlockingQueue<Object> queue, Blackhole bh) {
@@ -278,25 +273,5 @@ public void benchReactor(Blackhole bh) throws InterruptedException {
278273
latch.await();
279274
}
280275

281-
@Benchmark
282-
@OperationsPerInvocation(BATCH_SIZE)
283-
public void benchVertx(Blackhole bh) throws InterruptedException {
284-
CountDownLatch latch = new CountDownLatch(BATCH_SIZE);
285-
io.vertx.core.eventbus.MessageConsumer<Object> consumer = vertx.eventBus().localConsumer("benchmark.address", msg -> {
286-
bh.consume(msg.body());
287-
latch.countDown();
288-
});
289276

290-
CountDownLatch regLatch = new CountDownLatch(1);
291-
consumer.completionHandler(res -> regLatch.countDown());
292-
regLatch.await();
293-
294-
runProducers(() -> {
295-
for (int i = 0; i < EVENTS_PER_PRODUCER; i++) {
296-
vertx.eventBus().send("benchmark.address", EVENT);
297-
}
298-
});
299-
latch.await();
300-
consumer.unregister();
301-
}
302277
}

‎zthread-benchmark/src/main/java/io/github/namanoncode/zthread/benchmark/throughput/MpscEventBenchmark.java‎

Lines changed: 0 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -278,26 +278,4 @@ public void benchReactor(Blackhole bh) throws InterruptedException {
278278
latch.await();
279279
}
280280

281-
@Benchmark
282-
@OperationsPerInvocation(BATCH_SIZE)
283-
public void benchVertx(Blackhole bh) throws InterruptedException {
284-
CountDownLatch latch = new CountDownLatch(BATCH_SIZE);
285-
io.vertx.core.eventbus.MessageConsumer<Object> consumer = vertx.eventBus().localConsumer("benchmark.address", msg -> {
286-
bh.consume(msg.body());
287-
latch.countDown();
288-
});
289-
290-
// Wait for consumer to be registered to avoid dropping messages
291-
CountDownLatch regLatch = new CountDownLatch(1);
292-
consumer.completionHandler(res -> regLatch.countDown());
293-
regLatch.await();
294-
295-
runProducers(() -> {
296-
for (int i = 0; i < EVENTS_PER_PRODUCER; i++) {
297-
vertx.eventBus().send("benchmark.address", EVENT);
298-
}
299-
});
300-
latch.await();
301-
consumer.unregister();
302-
}
303281
}

‎zthread-benchmark/src/main/java/io/github/namanoncode/zthread/benchmark/throughput/SpscEventBenchmark.java‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -205,7 +205,9 @@ public void benchReactor(Blackhole bh) throws InterruptedException {
205205
@OperationsPerInvocation(BATCH_SIZE)
206206
public void benchVertx(Blackhole bh) throws InterruptedException {
207207
CountDownLatch latch = new CountDownLatch(BATCH_SIZE);
208-
io.vertx.core.eventbus.MessageConsumer<Object> consumer = vertx.eventBus().localConsumer("benchmark.address", msg -> {
208+
io.vertx.core.eventbus.MessageConsumer<Object> consumer = vertx.eventBus().localConsumer("benchmark.address");
209+
consumer.setMaxBufferedMessages(10_000_000);
210+
consumer.handler(msg -> {
209211
bh.consume(msg.body());
210212
latch.countDown();
211213
});
@@ -219,6 +221,6 @@ public void benchVertx(Blackhole bh) throws InterruptedException {
219221
vertx.eventBus().send("benchmark.address", "bench");
220222
}
221223
latch.await();
222-
consumer.unregister();
224+
consumer.unregister().toCompletionStage().toCompletableFuture().join();
223225
}
224226
}

0 commit comments

Comments
 (0)