Created
May 21, 2020 08:20
-
-
Save danielkec/828fd2afcc0841cb508a4b73f94e25b1 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package io.helidon.reactive.jmh; | |
| import java.util.List; | |
| import java.util.concurrent.Flow; | |
| import java.util.stream.Collectors; | |
| import java.util.stream.IntStream; | |
| import io.helidon.common.reactive.BufferedEmittingPublisher; | |
| import io.helidon.common.reactive.Multi; | |
| import io.helidon.common.reactive.OriginThreadPublisher; | |
| import org.openjdk.jmh.annotations.Benchmark; | |
| import org.openjdk.jmh.infra.Blackhole; | |
| import org.openjdk.jmh.runner.Runner; | |
| import org.openjdk.jmh.runner.options.Options; | |
| import org.openjdk.jmh.runner.options.OptionsBuilder; | |
| public class EmitterTest { | |
| public static void main(String[] args) throws Exception { | |
| Options opt = new OptionsBuilder() | |
| .include(EmitterTest.class.getSimpleName()) | |
| .forks(1) | |
| .build(); | |
| new Runner(opt).run(); | |
| } | |
| static final List<String> TEST_DATA = IntStream.range(0, 1_000_000) | |
| .mapToObj(i -> String.format("%d%d%d%d%d", i, i, i, i, i)) | |
| .collect(Collectors.toList()); | |
| @Benchmark | |
| public void testEmitterUnbounded(Blackhole bh) { | |
| BufferedEmittingPublisher<String> emitter = BufferedEmittingPublisher.create(); | |
| Multi.from(emitter).forEach(bh::consume); | |
| TEST_DATA.forEach(emitter::emit); | |
| } | |
| @Benchmark | |
| public void testEmitterReqMillion(Blackhole bh) { | |
| BufferedEmittingPublisher<String> emitter = BufferedEmittingPublisher.create(); | |
| TestSubscriber<String> subscriber = new TestSubscriber<>(); | |
| Multi.from(emitter).subscribe(subscriber); | |
| subscriber.request(1_000_000); | |
| TEST_DATA.forEach(emitter::emit); | |
| } | |
| @Benchmark | |
| public void testEmitterReqOneByOne(Blackhole bh) { | |
| BufferedEmittingPublisher<String> emitter = BufferedEmittingPublisher.create(); | |
| TestSubscriber<String> subscriber = new TestSubscriber<>(); | |
| Multi.from(emitter).subscribe(subscriber); | |
| for (String TEST_DATUM : TEST_DATA) { | |
| subscriber.request(1); | |
| emitter.emit(TEST_DATUM); | |
| } | |
| } | |
| @Benchmark | |
| public void testOTPUnbounded(Blackhole bh) { | |
| OriginThreadPublisher<String, String> emitter = new OriginThreadPublisher<>() {}; | |
| Multi.from(emitter).forEach(bh::consume); | |
| TEST_DATA.forEach(emitter::submit); | |
| } | |
| @Benchmark | |
| public void testOTPReqMillion(Blackhole bh) { | |
| OriginThreadPublisher<String, String> emitter = new OriginThreadPublisher<>() {}; | |
| TestSubscriber<String> subscriber = new TestSubscriber<>(); | |
| Multi.from(emitter).subscribe(subscriber); | |
| subscriber.request(1_000_000); | |
| TEST_DATA.forEach(emitter::submit); | |
| } | |
| @Benchmark | |
| public void testOTPReqOneByOne(Blackhole bh) { | |
| OriginThreadPublisher<String, String> emitter = new OriginThreadPublisher<>() {}; | |
| TestSubscriber<String> subscriber = new TestSubscriber<>(); | |
| Multi.from(emitter).subscribe(subscriber); | |
| TEST_DATA.forEach(data -> { | |
| subscriber.request(1); | |
| emitter.submit(data); | |
| }); | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment