[Benchmark] 使用ParallelStream时,CompletableFuture.get方法性能太差
- Dominant language
- Java
- Stars
- 41.6k
- Forks
- 26.4k
- Avg merge
- 15h 13m
- Merged PRs (30d)
- 4
Description
- [x] I have searched the [issues](https://github.com/apache/dubbo/issues) of this repository and believe that this is not a duplicate.
- [x] I have checked the [FAQ](https://github.com/apache/dubbo/blob/master/FAQ.md) of this repository and believe that this is not a duplicate.
### Environment
* Dubbo version: 2.7
* Operating System version: Mac OS 10.15
* Java version: 1.8
### Steps to reproduce this issue
before 2.7, DefaultFuture.get(time) like this code
```java
public static class ConsumerFuture {
private final Lock lock = new ReentrantLock();
private final Condition done = lock.newCondition();
private volatile T obj;
// get
public T get(int timeout) throws RemotingException {
if (timeout <= 0) {
timeout = Constants.DEFAULT_TIMEOUT;
}
if (!isDone()) {
long start = System.currentTimeMillis();
lock.lock();
try {
while (!isDone()) {
done.await(timeout, TimeUnit.MILLISECONDS);
if (isDone() || System.currentTimeMillis() - start >= timeout) {
break;
}
}
} catch (InterruptedException e) {
throw new RuntimeException(e);
} finally {
lock.unlock();
}
if (!isDone()) {
throw new TimeoutException(true, null, null);
}
}
return obj;
}
public boolean isDone() {
return obj != null;
}
// 结束值
private void doReceived(T data) {
lock.lock();
try {
this.obj = data;
if (done != null) {
done.signal();
}
} finally {
lock.unlock();
}
}
}
```
after 2.7, the code is `CompletableFuture.get`.
now I have do some benchmark test like this.
```java
// 评估吞吐率
@BenchmarkMode({Mode.Throughput})
// 执行1轮
@Fork(1)
// 每次执行使用8个线程
@Threads(8)
// 预热2次,每次预热10s
@Warmup(iterations = 2, time = 10)
// 执行6次测量,每次执行10s,共1分钟
@Measurement(iterations = 6, time = 10)
// 输出计算单位minutes
@OutputTimeUnit(TimeUnit.MINUTES)
// 共享变量,每次benchmark共享1个变量。
@State(Scope.Benchmark)
public class DefaultFutureBenchmark {
private ExecutorService pool = Executors.newFixedThreadPool(4, new NamedThreadFactory("Benchmark", true));
private void parallelExec(Runnable runnable) {
List> circle = new ArrayList<>();
for (int i = 0; i < 20; i++) {
ArrayList item = new ArrayList<>();
for (int j = 0; j < 25; j++) {
item.add(j);
}
circle.add(item);
}
// 此处测试模拟调用,共执行3次测试,具体3次代码在后文展示,
runTask(circle, runnable);
}
// 测试1, 两层parallelStream循环
private void runTask(List> circle, Runnable runnable) {
List> list = circle.parallelStream().map(item -> item.parallelStream().map(i -> {
runnable.run();
return "Success";
}).collect(Collectors.toList())
).collect(Collectors.toList());
assert list.size() == 20;
}
// 测试2, 仅内层使用parallelStream循环,相当于1层parallelStream
private void runTask(List> circle, Runnable runnable) {
List> list = circle.stream().map(item -> item.parallelStream().map(i -> {
runnable.run();
return "Success";
}).collect(Collectors.toList())
).collect(Collectors.toList());
assert list.size() == 20;
}
// 测试3, 单线程两次stream循环调用
private void runTask(List> circle, Runnable runnable) {
List> list = circle.stream().map(item -> item.stream().map(i -> {
runnable.run();
return "Success";
}).collect(Collectors.toList())
).collect(Collectors.toList());
assert list.size() == 20;
}
@Benchmark
public void testCompletableFutureGetTimeout() {
parallelExec(() -> {
try {
CompletableFuture future = new CompletableFuture<>();
pool.submit(() -> {
future.complete("Success");
});
future.get(100_000, TimeUnit.MILLISECONDS);
} catch (Exception e) {
// do nothing
}
});
}
@Benchmark
public void testCompletableFutureGet() {
parallelExec(() -> {
try {
CompletableFuture future = new CompletableFuture<>();
pool.submit(() -> {
future.complete("Success");
});
future.get();
} catch (Exception e) {
// do nothing
}
});
}
@Benchmark
public void testConsumerFuture() {
parallelExec(() -> {
try {
ConsumerFuture future = new ConsumerFuture<>();
pool.submit(() -> {
future.doReceived("Success");
});
future.get(10_0000);
} catch (Exception e) {
// do nothing
}
});
}
}
```
Result:
Test1:
```txt
Benchmark Mode Cnt Score Error Units
DefaultFutureBenchmark.testCompletableFutureGet thrpt 6 1489.322 ± 269.256 ops/min
DefaultFutureBenchmark.testCompletableFutureGetTimeout thrpt 6 1774.805 ± 334.626 ops/min
DefaultFutureBenchmark.testConsumerFuture thrpt 6 50869.454 ± 7393.916 ops/min
```
Test2:
```txt
Benchmark Mode Cnt Score Error Units
DefaultFutureBenchmark.testCompletableFutureGet thrpt 6 9611.754 ± 836.907 ops/min
DefaultFutureBenchmark.testCompletableFutureGetTimeout thrpt 6 8469.854 ± 1424.671 ops/min
DefaultFutureBenchmark.testConsumerFuture thrpt 6 54473.595 ± 5832.127 ops/min
```
Test3:
```txt
Benchmark Mode Cnt Score Error Units
DefaultFutureBenchmark.testCompletableFutureGet thrpt 6 93994.076 ± 4341.982 ops/min
DefaultFutureBenchmark.testCompletableFutureGetTimeout thrpt 6 47404.775 ± 11263.947 ops/min
DefaultFutureBenchmark.testConsumerFuture thrpt 6 50488.342 ± 3728.483 ops/min
```
Problem:
1. Parallelstream can seriously affect CompletableFuture.get Performance.(parallelStream会严重影响CompletableFuture.get的性能);
2. The influence of double-layer parallel stream is much greater than that of single-layer parallel stream.(双层parallelStream的影响比单层parallelStream的影响要大的多);
3. The performance of ConsumerFuture is not affected by parallel stream, and its performance is similar.(ConsumerFuture性能不因parallelStream的影响,其表现相差不大)
Contributor guide
Assessment
This issue has not been assessed yet.