Skip to content

Instantly share code, notes, and snippets.

@danielkec
Created May 21, 2020 08:20
Show Gist options
  • Select an option

  • Save danielkec/828fd2afcc0841cb508a4b73f94e25b1 to your computer and use it in GitHub Desktop.

Select an option

Save danielkec/828fd2afcc0841cb508a4b73f94e25b1 to your computer and use it in GitHub Desktop.
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