/
githubmirror
/
RxJava
Обзор
Документация
Войти
/
githubmirror
/
RxJava
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
4.x
src/jmh/java/io/reactivex/rxjava4/parallel/ParallelPerf.java
103 строки
3 KB
David Karnok
4.x: Streamable + takeUntil + groupBy + refactor + helpers (#8215)
05 июл 2026, 01:45
Не верифицирован
05 июл 2026, 01:45
580b39c
Код
Авторство
О чём код?
/* * Copyright (c) 2016-present, RxJava Contributors. * * 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.rxjava4.parallel; import java.util.Arrays; import java.util.concurrent.Flow.Publisher; import java.util.concurrent.TimeUnit; import org.openjdk.jmh.annotations.*; import org.openjdk.jmh.infra.Blackhole; import io.reactivex.rxjava4.core.*; import io.reactivex.rxjava4.core.config.StandardConcurrentBufferedConfig; import io.reactivex.rxjava4.functions.Function; import io.reactivex.rxjava4.schedulers.Schedulers; @SuppressWarnings("exports") @BenchmarkMode(Mode.Throughput) @Warmup(iterations = 5) @Measurement(iterations = 5, time = 1, timeUnit = TimeUnit.SECONDS) @Fork(value = 1, jvmArgsAppend = { "-XX:MaxInlineLevel=20" }) @OutputTimeUnit(TimeUnit.SECONDS) @State(Scope.Thread) public class ParallelPerf implements Function<Integer, Integer> { @Param({"10000"}) public int count; @Param({"1", "10", "100", "1000", "10000"}) public int compute; @Param({"1", "2", "3", "4"}) public int parallelism; Flowable<Integer> flatMap; Flowable<Integer> groupBy; Flowable<Integer> parallel; @Override public Integer apply(Integer t) { Blackhole.consumeCPU(compute); return t; } @Setup public void setup() { final int cpu = parallelism; Integer[] ints = new Integer[count]; Arrays.fill(ints, 777); Flowable<Integer> source = Flowable.fromArray(ints); flatMap = source.flatMap((Function<Integer, Publisher<Integer>>) v -> Flowable.just(v).subscribeOn(Schedulers.computation()) .map(ParallelPerf.this), new StandardConcurrentBufferedConfig(cpu)); groupBy = source.groupBy(new Function<Integer, Integer>() { int i; @Override public Integer apply(Integer v) { return (i++) % cpu; } }) .flatMap((Function<GroupedFlowable<Integer, Integer>, Publisher<Integer>>) g -> g.observeOn(Schedulers.computation()).map(ParallelPerf.this)); parallel = source.parallel(cpu).runOn(Schedulers.computation()).map(this).sequential(); } void subscribe(Flowable<Integer> f, Blackhole bh) { PerfAsyncConsumer consumer = new PerfAsyncConsumer(bh); f.subscribe(consumer); consumer.await(count); } @Benchmark public void flatMap(Blackhole bh) { subscribe(flatMap, bh); } @Benchmark public void groupBy(Blackhole bh) { subscribe(groupBy, bh); } @Benchmark public void parallel(Blackhole bh) { subscribe(parallel, bh); } }