ahuang98 commented on code in PR #22669: URL: https://github.com/apache/kafka/pull/22669#discussion_r3606033459
########## jmh-benchmarks/src/main/java/org/apache/kafka/jmh/raft/ElectionBenchmarks.java: ########## @@ -0,0 +1,82 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.kafka.jmh.raft; + +import org.apache.kafka.raft.RaftClientBenchmarkContext; +import org.apache.kafka.raft.RaftClientTestContext; + +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.Warmup; + +import java.io.IOException; +import java.util.Optional; +import java.util.concurrent.TimeUnit; + +/** + * Benchmarks for the leader-election path. The outer class is intentionally not a JMH {@code @State}: + * each benchmark declares the starting state it needs as a nested {@code @State} parameter, so + * different election scenarios (e.g. a future Prospective or Candidate start) can have their own + * setup without forcing a single shared {@code @Setup} on the whole class. + */ +@BenchmarkMode(Mode.SingleShotTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +@Warmup(iterations = RaftClientBenchmarkContext.SINGLE_SHOT_WARMUP_ITERATIONS) +@Measurement(iterations = RaftClientBenchmarkContext.SINGLE_SHOT_MEASUREMENT_ITERATIONS) +@Fork(RaftClientBenchmarkContext.SINGLE_SHOT_FORKS) +public class ElectionBenchmarks { + + /** + * Starting state: the local node is Unattached in a {@code voterCount}-node cluster. A fresh + * context is built per invocation because driving the election to completion consumes it. Review Comment: I'm not sure if I understand the last sentence of this comment (e.g. what does 'consumes it' mean) ########## jmh-benchmarks/src/main/java/org/apache/kafka/jmh/raft/KRaftBenchmarkingCounters.java: ########## @@ -0,0 +1,156 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.kafka.jmh.raft; + +import org.apache.kafka.common.protocol.ApiKeys; +import org.apache.kafka.raft.RaftClientBenchmarkContext; + +import org.openjdk.jmh.annotations.AuxCounters; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.infra.BenchmarkParams; + +import java.util.Optional; + +/** + * Secondary, machine-independent work counters reported by the raft benchmarks alongside the timing + * score, as {@code benchmark:counter} rows. + * + * <p>Throughout this class, an <em>operation</em> is JMH's unit of work: a single invocation of a + * {@code @Benchmark}-annotated method. (One operation equals one invocation here because we don't use + * {@code @OperationsPerInvocation}.) JMH reports the timing score in {@code ns/op}, and these work + * counters are reported {@code PerOp} to match. + * + * <p>The per-operation values are integer-exact and should be stable across a correct refactor of + * {@code KafkaRaftClient}: a flush count moving from 1 to 2 per operation is a behavioral diff, not + * measurement noise. The counters that are zero on a path (e.g. log flushes on a caught-up fetch) + * are the most useful tripwires, since zero is speed-independent. + */ +@State(Scope.Thread) +@AuxCounters(AuxCounters.Type.EVENTS) +public class KRaftBenchmarkingCounters { + // Private accumulators: not reported directly (we report the per-op values below). Being private, + // JMH does not touch them between iterations, so reset() must zero them. + private long logFlushesTotal; + private long logReadsTotal; + private long logTruncationsTotal; + private long rpcRequestsSentTotal; + private long rpcResponsesSentTotal; + private long quorumStateWritesTotal; + private long quorumStateReadsTotal; + + // The number of operations (i.e. @Benchmark method invocations) measured in the iteration, and the + // divisor for the per-operation values below. + private long operations; + + // The divisor for the per-op methods below: (forks x measurement iterations). Set once per fork + // by captureRunShape(). + private double measurementDataPoints = 1.0; + + /** + * Captures the number of measurement data points — {@code forks x measurement iterations} — that + * JMH will SUM the {@code *PerOp()} methods over ({@code Type.EVENTS} secondary results are + * SUM-aggregated across iterations and forks). Each per-op method pre-divides by this count so that + * the SUM reports the exact per-operation value (e.g. {@code logReadsPerOp = 1.0}) in the summary + * row. Reading it from {@link BenchmarkParams} tracks the actual run shape (including + * {@code -f}/{@code -i} overrides) rather than hardcoding the annotation values. + */ + @Setup(Level.Trial) + public void captureRunShape(BenchmarkParams params) { + if (params.getThreads() != 1) { + throw new IllegalStateException( + "raft benchmarks are single-threaded (one client over shared mocks); got " + + params.getThreads() + " threads"); + } + // forks() is 0 when forking is disabled (in-process), which is still one set of iterations. + int forks = Math.max(1, params.getForks()); + measurementDataPoints = (double) forks * params.getMeasurement().getCount(); + } + + @Setup(Level.Iteration) + public void reset() { + logFlushesTotal = 0; + logReadsTotal = 0; + logTruncationsTotal = 0; + rpcRequestsSentTotal = 0; + rpcResponsesSentTotal = 0; + quorumStateWritesTotal = 0; + quorumStateReadsTotal = 0; + operations = 0; + } + + /** Review Comment: It was a bit hard for me to parse the first sentence of this method description (what are 'these counters'?) Maybe excluding some specifics might make it easier to understand generally what this method does? ``` Records this invocation's work deltas from `context`, drains the expected in-flight request/response RPCs, and counts one completed benchmark operation. ``` ########## jmh-benchmarks/src/main/java/org/apache/kafka/jmh/raft/LeaderBenchmarks.java: ########## @@ -0,0 +1,100 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.kafka.jmh.raft; + +import org.apache.kafka.common.protocol.ApiKeys; +import org.apache.kafka.raft.RaftClientBenchmarkContext; +import org.apache.kafka.raft.RaftClientTestContext; + +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Fork; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Measurement; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.Warmup; + +import java.util.Optional; +import java.util.concurrent.TimeUnit; + +/** + * Benchmarks for the leader request-handling path. The outer class is intentionally not a JMH + * {@code @State}: each benchmark declares the starting state it needs as a nested {@code @State} + * parameter, so future leader scenarios (e.g. a lagging-follower fetch or a commit) can have their own + * setup without forcing a single shared {@code @Setup} on the whole class. + */ +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +@Warmup(iterations = RaftClientBenchmarkContext.AVERAGE_TIME_WARMUP_ITERATIONS) +@Measurement(iterations = RaftClientBenchmarkContext.AVERAGE_TIME_MEASUREMENT_ITERATIONS) +@Fork(RaftClientBenchmarkContext.AVERAGE_TIME_FORKS) +public class LeaderBenchmarks { + + /** + * Starting state: the local node is leader with the high watermark at the log end. + */ + @State(Scope.Thread) + public static class LeaderWithHwmAtLogEnd { + static final int VOTER_COUNT = 3; + + RaftClientBenchmarkContext benchmark; + RaftClientTestContext context; + + int epoch; + long endOffset; + + @Setup(Level.Trial) + public void setup() throws Exception { + benchmark = RaftClientBenchmarkContext.leader( + VOTER_COUNT, + 0, + RaftClientBenchmarkContext.DEFAULT_KRAFT_VERSION, + RaftClientBenchmarkContext.DEFAULT_RAFT_PROTOCOL); + context = benchmark.testContext(); + context.advanceLocalLeaderHighWatermarkToLogEndOffset(); + epoch = context.currentEpoch(); + endOffset = benchmark.logEndOffset(); + benchmark.zeroCountersOnSetup(); + } + } + + /** + * Leader handles a FETCH from a fully caught-up follower (fetch offset == log end offset) that asks + * not to wait ({@code maxWaitMs = 0}), so the leader replies immediately rather than deferring. + * + * <p>Note: a real caught-up follower long-polls with {@code maxWaitMs > 0}, and such a fetch is Review Comment: what is the benefit of benchmarking this scenario if it is not expected to occur in reality? (why not benchmark the scenario where a follower is not fully caught up and the leader has enough new data to immediately respond?) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
