Jmh测试框架和flink benchmark工程

[toc]

本文着重介绍对jmh框架的理解和使用,以及在flink-benchmark中的应用。

概述

JMH是一个由OpenJDK/Oracle里面那群开发了Java编译器的大牛们所开发的Micro Benchmark Framework。何谓Micro Benchmark呢?简单地说就是在method层面上的 benchmark,精度可以精确到微秒级。可以看出JMH主要使用在当你已经找出了热点函数,而需要对热点函数进行进一步的优化时,就可以使用JMH对优化的效果进行定量的分析。在flink-benchmark框架中主要用来对network和state进行定量分析,来确保每次对这两个模块的修改不会导致性能的regress.

http://apache-flink-mailing-list-archive.1008284.n3.nabble.com/Codespeed-deployment-for-Flink-td24274.html

比较典型的使用场景有:

想定量地知道某个函数需要执行多长时间,以及执行时间和输入n的相关性,一个函数有两种不同实现(例如实现A使用了FixedThreadPool,实现B使用了ForkJoinPool),不知道哪种实现性能更好.

参数含义

测试case:

https://github.com/Aitozi/test-case/blob/68a6d5dbc28220cb54b8057d147632f3aed6bd31/src/main/java/jmh/SimpleBenchT.java

@BenchmarkMode

基准测试类型。这里选择的是Throughput也就是吞吐量。根据源码点进去,每种类型后面都有对应的解释,比较好理解,吞吐量会得到单位时间内可以进行的操作数。

  • Throughput: 整体吞吐量,例如“1秒内可以执行多少次调用”。
  • AverageTime: 调用的平均时间,例如“每次调用平均耗时xxx毫秒”。
  • SampleTime: 随机取样,最后输出取样结果的分布,例如“99%的调用在xxx毫秒以内,99.99%的调用在xxx毫秒以内”
  • SingleShotTime: 以上模式都是默认一次iteration是1s,唯有SingleShotTime是只运行一次。往往同时把warmup次数设为0,用于测试冷启动时的性能。
  • All 执行所有的类型测试

@Warmup

进行基准测试前需要进行预热。一般我们前几次进行程序测试的时候都会比较慢, 所以要让程序进行几轮预热,保证测试的准确性。其中的参数iterations也就非常好理解了,就是预热轮数。

为什么需要预热?因为JVM的JIT机制的存在,如果某个函数被调用多次之后,JVM会尝试将其编译成为机器码从而提高执行速度。所以为了让benchmark的结果更加接近真实情况就需要进行预热。

@Measurement

测试参数

  • iterations 进行测试的轮次
  • time 每轮进行的时长
  • timeUnit 时长单位

都是一些基本的参数,可以根据具体情况调整。一般比较重的东西可以进行大量的测试,放到服务器上运行。

@Threads

每个进程中的测试线程,这个非常好理解,根据具体情况选择,一般为cpu乘以2。

@Fork

进行fork的次数。如果fork数是2的话,则JMH会fork出两个进程来进行测试。

@OutputTimeUnit

这个比较简单了,基准测试结果的时间类型。一般选择秒、毫秒、微秒。

@Benchmark

方法级注解,表示该方法是需要进行benchmark的对象,用法和JUnit的@Test类似。

@Param

属性级注解,@Param可以用来指定某项参数的多种情况。特别适合用来测试一个函数在不同的参数输入的情况下的性能。

@Setup

方法级注解,这个注解的作用就是我们需要在测试之前进行一些准备工作,比如对一些数据的初始化之类的。

@TearDown

方法级注解,这个注解的作用就是我们需要在测试之后进行一些结束工作,比如关闭线程池,数据库连接等的,主要用于资源的回收等。

@State

当使用@Setup参数的时候,必须在类上加这个参数,不然会提示无法运行。

State用于声明某个类是一个“状态”,然后接受一个Scope参数用来表示该状态的共享范围。 因为很多benchmark会需要一些表示状态的类,JMH允许你把这些类以依赖注入的方式注入到 benchmark函数里。Scope主要分为三种。

Thread: 该状态为每个线程独享。
Group: 该状态为同一个组里面所有线程共享。
Benchmark: 该状态在所有线程间共享。

关于State的用法,官方的code sample里有比较好的例子。

其他

  • CompilerControl控制 compiler 的行为,例如强制 inline,不允许编译等。

  • Group 可以把多个 benchmark 定义为同一个 group,则它们会被同时执行,主要用于测试多个相互之间存在影响的方法。

  • Level 用于控制 @Setup,@TearDown 的调用时机,默认是 Level.Trial,即benchmark开始前和结束后。

  • Profiler JMH 支持一些 profiler,可以显示等待时间和运行时间比,热点函数等。

serialization

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
@Benchmark
@OperationsPerInvocation(value = RECORDS_PER_INVOCATION)
/**
* 告诉jmh该测试方法内部会执行RECORDS_PER_INVOCATION次数,最后计算score的时候需要
* 把这个考虑上.http://javadox.com/org.openjdk.jmh/jmh-core/0.9/org/openjdk/jmh/annotations/OperationsPerInvocation.html
/
public void serializerKryoWithoutRegistration() throws Exception {
LocalStreamEnvironment env =
StreamExecutionEnvironment.createLocalEnvironment(4);
env.getConfig().enableForceKryo();

env.addSource(new PojoSource(RECORDS_PER_INVOCATION, 10))
.rebalance()
.addSink(new DiscardingSink<>());

env.execute();
}

以及: SerializationFrameworkMiniBenchmarks

keyby

测试在tuple和array上的keyby性能

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
@Benchmark
@OperationsPerInvocation(value = KeyByBenchmarks.TUPLE_RECORDS_PER_INVOCATION)
public void tupleKeyBy() throws Exception {
LocalStreamEnvironment env =
StreamExecutionEnvironment.createLocalEnvironment(4);

env.addSource(new IncreasingTupleSource(TUPLE_RECORDS_PER_INVOCATION, 10))
.keyBy(0)
.addSink(new DiscardingSink<>());

env.execute();
}

@Benchmark
@OperationsPerInvocation(value = KeyByBenchmarks.ARRAY_RECORDS_PER_INVOCATION)
public void arrayKeyBy() throws Exception {
LocalStreamEnvironment env =
StreamExecutionEnvironment.createLocalEnvironment(4);

env.addSource(new IncreasingArraySource(ARRAY_RECORDS_PER_INVOCATION, 10))
.keyBy(0)
.addSink(new DiscardingSink<>());

env.execute();
}

state backend

1
2
3
4
5
6
 source
.map(new MultiplyIntLongByTwo())
.keyBy(record -> record.key)
.window(windowAssigner)
.reduce(new SumReduceIntLong())
.addSink(new CollectSink());

利用window窗口聚合,对不同backend做不同的基准测试,测试的和state api不够直接

window

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
@Benchmark
public void tumblingWindow(TimeWindowContext context) throws Exception {
IntLongApplications.reduceWithWindow(context.source, TumblingEventTimeWindows.of(Time.seconds(10_000)));
context.execute();
}

@Benchmark
public void slidingWindow(TimeWindowContext context) throws Exception {
IntLongApplications.reduceWithWindow(context.source, SlidingEventTimeWindows.of(Time.seconds(10_000), Time.seconds(1000)));
context.execute();
}

@Benchmark
public void sessionWindow(TimeWindowContext context) throws Exception {
IntLongApplications.reduceWithWindow(context.source, EventTimeSessionWindows.withGap(Time.seconds(500)));
context.execute();
}

network

测试吞吐和延迟:

  • StreamNetworkThroughputBenchmarkExecutor
  • StreamNetworkLatencyBenchmarkExecutor

state operations

https://github.com/dataArtisans/flink-benchmarks/pull/13

###参考

https://www.xncoding.com/2018/01/07/java/jmh.html
http://tutorials.jenkov.com/java-performance/jmh.html
https://www.cnkirito.moe/java-jmh/ 详细介绍平常测试代码中的陷阱

谢谢支持