This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12418-4c4fd615d677695b619ca0e04fcfd82e104cd4a9 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit d7e9931bea0976547e9730675e7019bbef88ad83 Author: Gangavarapu Vivek <[email protected]> AuthorDate: Thu Oct 1 04:59:34 2026 +0000 [Test][Zeta] Add a checkpoint scheduling delay benchmark (#12418) --- .github/workflows/benchmarks.yml | 1 + .../benchmark/CheckpointSchedulingBenchmark.java | 157 ++++++ .../benchmark/checkpoint/BenchmarkReflection.java | 48 ++ .../checkpoint/CheckpointSchedulingFixture.java | 607 +++++++++++++++++++++ .../checkpoint/FakeTaskCheckpointManager.java | 160 ++++++ .../CheckpointSchedulingBenchmarkTest.java | 42 ++ .../CheckpointSchedulingFixtureTest.java | 150 +++++ tools/benchmarks/save_jmh_result.py | 21 +- tools/benchmarks/test_save_jmh_result.py | 26 + 9 files changed, 1211 insertions(+), 1 deletion(-) diff --git a/.github/workflows/benchmarks.yml b/.github/workflows/benchmarks.yml index 62721a1476..3d31440bd8 100644 --- a/.github/workflows/benchmarks.yml +++ b/.github/workflows/benchmarks.yml @@ -40,6 +40,7 @@ on: - 'ProtoStuffSerializerBenchmark' - 'SeaTunnelPipelineBenchmark' - 'CheckpointingTimeBenchmark' + - 'CheckpointSchedulingBenchmark' - 'CheckpointStorageBenchmark' - 'IMapJobStorageBenchmark' - 'IMapDagStorageBenchmark' diff --git a/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/CheckpointSchedulingBenchmark.java b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/CheckpointSchedulingBenchmark.java new file mode 100644 index 0000000000..8e2b855310 --- /dev/null +++ b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/CheckpointSchedulingBenchmark.java @@ -0,0 +1,157 @@ +/* + * 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.seatunnel.benchmark; + +import org.apache.seatunnel.benchmark.checkpoint.CheckpointSchedulingFixture; + +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.TearDown; +import org.openjdk.jmh.annotations.Threads; +import org.openjdk.jmh.annotations.Warmup; +import org.openjdk.jmh.runner.Runner; +import org.openjdk.jmh.runner.RunnerException; +import org.openjdk.jmh.runner.options.Options; +import org.openjdk.jmh.runner.options.OptionsBuilder; +import org.openjdk.jmh.runner.options.VerboseMode; + +import java.util.concurrent.TimeUnit; + +/** + * Measures how late a periodic checkpoint trigger runs after it is due, on a real member with + * {@code pipelineNum} real checkpoint coordinators. + * + * <p>This is the part of checkpointing the scheduling model decides. Checkpoint completion time, + * which {@link CheckpointingTimeBenchmark} measures, is dominated by the barrier round-trip and is + * close to blind to how the trigger was scheduled. + * + * <p>Each invocation measures one trigger. The untimed setup picks the coordinator due soonest and + * returns exactly when its trigger is due; the timed body spins until that trigger has created its + * pending checkpoint. The score is therefore the scheduling delay plus the trigger's own decision + * logic up to creating the checkpoint, and nothing else. {@code Level.Invocation} is normally + * discouraged, but here the delay being measured is in microseconds while the setup overhead JMH + * leaves outside the timed region is tens of nanoseconds. See {@link CheckpointSchedulingFixture} + * for how due times are known and which triggers are skipped. + * + * <p>The checkpoint interval is {@code pipelineNum * triggerSpacingMillis}, so every point of the + * sweep sees the same rate of due triggers and the same checkpoint load on storage, and only the + * number of coordinators changes. It is floored at {@link #MIN_MEASURABLE_INTERVAL_MILLIS}, so the + * smallest pipeline counts sample less often. Each job has one pipeline: many jobs is what "many + * active pipelines" means on a member, and it keeps per-job checkpoint state from becoming the + * bottleneck. + */ +@BenchmarkMode(Mode.SampleTime) +@OutputTimeUnit(TimeUnit.MICROSECONDS) +@Threads(1) +@Warmup(iterations = 2, time = 10) +@Measurement(iterations = 5, time = 10) +@Fork( + value = 3, + jvmArgsAppend = { + "-Xms4g", + "-Xmx4g", + "-XX:+UseG1GC", + "-XX:+AlwaysPreTouch", + "-XX:+DisableExplicitGC", + "-XX:ActiveProcessorCount=4", + "-Djava.net.preferIPv4Stack=true" + }) +public class CheckpointSchedulingBenchmark extends BenchmarkBase { + + /** + * Floor for the checkpoint interval. A checkpoint takes a few milliseconds here, and an + * interval close to that sends triggers down the pending re-arm path instead of measuring them, + * which is what a 20 ms interval at one pipeline would do. This is a floor for measuring, above + * the lowest interval SeaTunnel accepts, which the fixture enforces separately. + */ + static final long MIN_MEASURABLE_INTERVAL_MILLIS = 200L; + + public static void main(String[] args) throws RunnerException { + Options options = + new OptionsBuilder() + .verbosity(VerboseMode.NORMAL) + .include( + ".*" + + CheckpointSchedulingBenchmark.class.getCanonicalName() + + ".*") + .build(); + + new Runner(options).run(); + } + + @Benchmark + public long periodicTriggerDelay(CoordinatorsState state) { + return state.fixture.awaitTrigger(); + } + + @State(Scope.Thread) + public static class CoordinatorsState { + + @Param({"1", "10", "100", "500"}) + private int pipelineNum; + + @Param({"20"}) + private long triggerSpacingMillis; + + private CheckpointSchedulingFixture fixture; + + @Setup(Level.Trial) + public void setUp() throws Exception { + fixture = + new CheckpointSchedulingFixture( + pipelineNum, + Math.max( + MIN_MEASURABLE_INTERVAL_MILLIS, + pipelineNum * triggerSpacingMillis)); + fixture.setUp(); + System.out.printf( + "# checkpoint scheduler threads for %d pipelines: %d; measuring %d of them%n", + pipelineNum, fixture.countSchedulerThreads(), fixture.probeCount()); + } + + @Setup(Level.Iteration) + public void setUpIteration() { + fixture.beginIteration(); + } + + @Setup(Level.Invocation) + public void awaitDueTrigger() { + fixture.awaitNextDueTrigger(); + } + + @TearDown(Level.Iteration) + public void tearDownIteration() { + System.out.println("# " + fixture.iterationReport()); + fixture.endIteration(); + } + + @TearDown(Level.Trial) + public void tearDown() throws Exception { + fixture.tearDown(); + } + } +} diff --git a/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/BenchmarkReflection.java b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/BenchmarkReflection.java new file mode 100644 index 0000000000..f970b54e59 --- /dev/null +++ b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/BenchmarkReflection.java @@ -0,0 +1,48 @@ +/* + * 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.seatunnel.benchmark.checkpoint; + +import java.lang.reflect.Field; + +/** Read access to engine internals that a benchmark fixture observes but has no getter for. */ +final class BenchmarkReflection { + + private BenchmarkReflection() {} + + /** + * Returns the named declared field, made accessible. + * + * @throws IllegalStateException naming the class and field when the field no longer exists, so + * an engine-internals rename fails the benchmark at setup instead of skewing its numbers + */ + static Field requireField(Class<?> owner, String name) { + try { + Field field = owner.getDeclaredField(name); + field.setAccessible(true); + return field; + } catch (NoSuchFieldException e) { + throw new IllegalStateException( + "Benchmark fixture reads " + + owner.getName() + + "#" + + name + + ", which no longer exists; update the fixture to the engine change", + e); + } + } +} diff --git a/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/CheckpointSchedulingFixture.java b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/CheckpointSchedulingFixture.java new file mode 100644 index 0000000000..2dc7cccb54 --- /dev/null +++ b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/CheckpointSchedulingFixture.java @@ -0,0 +1,607 @@ +/* + * 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.seatunnel.benchmark.checkpoint; + +import org.apache.seatunnel.benchmark.storage.SeaTunnelStorageEnvironmentContext; +import org.apache.seatunnel.engine.common.Constant; +import org.apache.seatunnel.engine.common.config.EngineConfig; +import org.apache.seatunnel.engine.common.config.SeaTunnelConfig; +import org.apache.seatunnel.engine.common.config.server.CheckpointConfig; +import org.apache.seatunnel.engine.server.SeaTunnelServer; +import org.apache.seatunnel.engine.server.checkpoint.CheckpointCoordinator; +import org.apache.seatunnel.engine.server.checkpoint.CheckpointPlan; +import org.apache.seatunnel.engine.server.execution.TaskGroupLocation; +import org.apache.seatunnel.engine.server.execution.TaskLocation; + +import com.hazelcast.map.IMap; +import com.hazelcast.spi.impl.NodeEngine; + +import java.lang.reflect.Field; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.SynchronousQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.LockSupport; +import java.util.stream.Collectors; + +/** + * One real SeaTunnel member running {@code pipelineNum} real checkpoint coordinators, one per job, + * whose periodic triggers the benchmark observes one at a time. + * + * <p>The coordinators are production code end to end; only the tasks are fakes, see {@link + * FakeTaskCheckpointManager}. Whichever checkpoint scheduler the engine on the classpath uses is + * the one being measured, so the same fixture runs unchanged on every revision being compared. + * + * <p>A trigger is observed through the coordinator's {@code pendingCounter}, which goes from 0 to 1 + * when the trigger body creates a pending checkpoint. The time a trigger is due is derived from the + * previous observation plus one interval, since the coordinator re-arms itself with that delay. A + * coordinator whose previous trigger was not observed has no known phase: it is taken out of the + * rotation, counted as a skip, and resynchronised later. + * + * <p>Only every {@code probeStride}-th coordinator is measured; the rest are load. The probes start + * at least {@link #MIN_PROBE_SPACING_NANOS} apart, which keeps two probes from coming due within + * one sample of each other even after their phases drift. Measuring every coordinator instead would + * skip whichever trigger came due while another was being measured, and those are the triggers that + * bunch up under contention, so the skips would bias the result towards short delays exactly where + * the scheduler is under the most pressure. Which coordinators are probes is fixed by index before + * anything is measured, so probe samples carry no such selection, and probe triggers still collide + * freely with the load. + */ +public final class CheckpointSchedulingFixture { + + /** Covers the per-pipeline pools and the member-wide scheduler threads alike. */ + static final String SCHEDULER_THREAD_NAME_PREFIX = "checkpoint-"; + + /** + * The coordinator has no getter for its count of in-flight checkpoints. It is read, never + * written, to see when a trigger has created a pending checkpoint. The trigger body increments + * it just before re-arming the next trigger, which is what makes "observed + interval" the next + * due time. If the field is renamed or removed, loading this class fails naming it. + */ + private static final Field PENDING_COUNTER_FIELD = + BenchmarkReflection.requireField(CheckpointCoordinator.class, "pendingCounter"); + + /** Lowest interval {@code CheckpointConfig} accepts. */ + static final long MIN_CHECKPOINT_INTERVAL_MILLIS = CheckpointConfig.MINIMAL_CHECKPOINT_TIME; + + /** Share of due triggers that may be skipped before an iteration is rejected. */ + static final double MAX_SKIP_RATIO = 0.1; + + /** + * Fewest due triggers an iteration needs before {@link #MAX_SKIP_RATIO} is enforced. A short + * smoke iteration, such as the CI benchmark job's one second, sees a handful of due triggers, + * where a single skip is already above the ratio and says nothing about the run. + */ + static final long MIN_DUE_TRIGGERS_FOR_SKIP_CHECK = 50; + + private static final int PIPELINE_ID = 1; + private static final long FIRST_JOB_ID = 1_000L; + private static final long NOT_SYNCED = Long.MIN_VALUE; + + /** + * Bounds of how long before a due trigger the setup stops parking and starts spinning. Parking + * alone would add the OS timer slack to the start of the measured window, and parking wakes + * late by that slack: tens of microseconds on Linux, about 5 ms on macOS. The window starts at + * the minimum and widens to the largest overshoot seen plus the minimum, up to the maximum. + */ + private static final long MIN_SPIN_WINDOW_NANOS = TimeUnit.MILLISECONDS.toNanos(1); + + private static final long MAX_SPIN_WINDOW_NANOS = TimeUnit.MILLISECONDS.toNanos(10); + + /** + * Resync waits at most this many intervals, plus {@link #PENDING_REARM_NANOS}, for every + * coordinator to trigger once. + */ + private static final int RESYNC_INTERVALS = 3; + + /** + * A trigger that finds a checkpoint still pending re-arms itself after 500 ms instead of one + * interval; resync allows for one such re-arm, with margin. + */ + private static final long PENDING_REARM_NANOS = TimeUnit.SECONDS.toNanos(1); + + /** Least spacing between the start phases of two measured coordinators. */ + private static final long MIN_PROBE_SPACING_NANOS = TimeUnit.MILLISECONDS.toNanos(100); + + private static final long EXECUTOR_KEEP_ALIVE_SECONDS = 60L; + private static final long SHUTDOWN_TIMEOUT_SECONDS = 30L; + + private final int pipelineNum; + private final long intervalMillis; + private final long intervalNanos; + private int probeStride; + private Set<Thread> preexistingSchedulerThreads = Collections.emptySet(); + + private final AtomicReference<Throwable> failure = new AtomicReference<>(); + private final List<FakeTaskCheckpointManager> managers = new ArrayList<>(); + + private SeaTunnelStorageEnvironmentContext environment; + private ThreadPoolExecutor coordinatorExecutor; + private AtomicInteger[] pendingCounters; + private long[] lastTriggerNanos; + + private long spinWindowNanos = MIN_SPIN_WINDOW_NANOS; + private int current = -1; + private long sampled; + private long skippedPending; + private long skippedCollided; + private long skippedOverrun; + private long skippedEarly; + + /** + * @param pipelineNum number of jobs, each with one single-task pipeline + * @param intervalMillis checkpoint interval of every job + */ + public CheckpointSchedulingFixture(int pipelineNum, long intervalMillis) { + this.pipelineNum = pipelineNum; + this.intervalMillis = intervalMillis; + this.intervalNanos = TimeUnit.MILLISECONDS.toNanos(intervalMillis); + } + + /** + * Starts the member, creates the coordinators and starts their pipelines staggered across one + * interval, so due triggers arrive as a steady stream rather than a burst. Returns once every + * coordinator's trigger phase is known. + */ + public void setUp() throws Exception { + validateParameters(); + probeStride = probeStride(pipelineNum, intervalNanos); + preexistingSchedulerThreads = new HashSet<>(liveSchedulerThreads()); + environment = new SchedulingEnvironmentContext(); + environment.setUp(); + SeaTunnelServer server = environment.getServer(); + EngineConfig engineConfig = server.getSeaTunnelConfig().getEngineConfig(); + coordinatorExecutor = createCoordinatorExecutor(engineConfig); + createManagers(server, engineConfig.getCheckpointConfig()); + startStaggered(); + // Every coordinator, not only the probes: once each has triggered, each has armed its + // scheduler, so countSchedulerThreads() sees the full thread cost of pipelineNum pipelines. + resync(1); + } + + /** Resynchronises the coordinators that were taken out of the rotation, then resets counts. */ + public void beginIteration() { + resync(probeStride); + sampled = 0; + skippedPending = 0; + skippedCollided = 0; + skippedOverrun = 0; + skippedEarly = 0; + } + + /** + * Picks the coordinator due soonest and returns exactly when its trigger is due, with its + * previous checkpoint completed and the trigger not yet run. Not measured. + */ + public void awaitNextDueTrigger() { + checkFailure(); + while (true) { + int next = soonestSynced(); + long expected = lastTriggerNanos[next] + intervalNanos; + long parkTarget = expected - spinWindowNanos; + long start = System.nanoTime(); + if (start >= expected) { + // It came due while the previous sample was being taken; it may have run + // unobserved, so its phase is lost. + skipCollided(next); + continue; + } + // Close enough to the due time already (the previous sample ended late in this + // trigger's window): skip parking and go straight to the spin. + if (start < parkTarget) { + parkUntil(parkTarget); + widenSpinWindow(System.nanoTime() - parkTarget); + } + // Counter first, then the clock: if the clock is still before the due time, the + // counter was read before the trigger could have run. + boolean pending = pendingCounters[next].get() != 0; + if (System.nanoTime() >= expected) { + // Parking overshot the due time; the trigger may have run unobserved. + skipOverrun(next); + continue; + } + if (pending) { + // A checkpoint is pending before the trigger is due. Usually the previous one is + // still running and this trigger takes the pending re-arm path; it can also be this + // trigger having run early against an estimate made from a late observation, when + // the measuring thread was descheduled. Either way the sample is unusable. + skipPending(next); + continue; + } + while (System.nanoTime() < expected) { + // Spin through the last stretch so the measured window starts on time. + } + if (pendingCounters[next].get() != 0) { + // The counter read 0 before the due time and only a trigger raises it, so the + // trigger ran no later than its estimated due time: the estimate was late. + skipEarly(next); + continue; + } + current = next; + return; + } + } + + /** + * Spins until the due coordinator's trigger has created its pending checkpoint. This is the + * measured part: it starts when the trigger is due and ends when the trigger has run. + * + * @return the time the trigger was observed + */ + public long awaitTrigger() { + AtomicInteger pendingCounter = pendingCounters[current]; + long deadline = System.nanoTime() + intervalNanos; + // The window starts slightly before the real deadline, and that is intended. The trigger + // body increments pendingCounter just before it re-arms the next trigger with + // schedule(..., interval), so "observed + interval" is a few microseconds early and each + // sample also includes the tail of the previous trigger body. That code is identical on + // every scheduler being compared, so the offset is equal on both sides and cancels in a + // comparison. Do not "fix" it by moving the start later: there is no observable point + // closer to the real deadline without changing engine code. + while (pendingCounter.get() == 0) { + if (System.nanoTime() > deadline) { + throw new IllegalStateException( + "Checkpoint trigger of job " + + (FIRST_JOB_ID + current) + + " did not run within one interval of being due"); + } + } + long observed = System.nanoTime(); + lastTriggerNanos[current] = observed; + sampled++; + return observed; + } + + /** + * Summarises the iteration's skips. Printed for every iteration, passing or not, so a run close + * to the limit stays visible. + */ + public String iterationReport() { + return String.format( + "measured %d of %d due triggers; skipped %d with a checkpoint already pending " + + "before the due time, %d that came due while another trigger was being " + + "measured, %d where parking overran the due time, %d that ran before " + + "their estimated due time", + sampled, + dueTriggers(), + skippedPending, + skippedCollided, + skippedOverrun, + skippedEarly); + } + + /** + * Rejects the iteration if more than {@link #MAX_SKIP_RATIO} of due triggers could not be + * measured, so a run where most triggers took the pending re-arm path produces an error rather + * than a number. Enforced once at least {@link #MIN_DUE_TRIGGERS_FOR_SKIP_CHECK} triggers were + * due; shorter iterations still print their counts. + */ + public void endIteration() { + checkFailure(); + if (isSkipShareTooHigh(dueTriggers(), dueTriggers() - sampled)) { + throw new IllegalStateException( + String.format( + "%d pipelines: %s. That is above the %.0f%% limit, so the measured " + + "delays would not be representative", + pipelineNum, iterationReport(), MAX_SKIP_RATIO * 100)); + } + } + + static boolean isSkipShareTooHigh(long due, long skipped) { + return due >= MIN_DUE_TRIGGERS_FOR_SKIP_CHECK && skipped > MAX_SKIP_RATIO * due; + } + + private long dueTriggers() { + return sampled + skippedPending + skippedCollided + skippedOverrun + skippedEarly; + } + + public void tearDown() throws Exception { + try { + // Executor first: a coordinator still finishing its asynchronous start would otherwise + // arm its first trigger on the scheduler that cancelling just replaced, and that + // thread would outlive the fixture. + if (coordinatorExecutor != null) { + coordinatorExecutor.shutdownNow(); + if (!coordinatorExecutor.awaitTermination( + SHUTDOWN_TIMEOUT_SECONDS, TimeUnit.SECONDS)) { + throw new IllegalStateException("Coordinator executor did not stop"); + } + } + for (FakeTaskCheckpointManager manager : managers) { + manager.cancelCheckpoint(PIPELINE_ID); + } + managers.clear(); + } finally { + coordinatorExecutor = null; + if (environment != null) { + environment.tearDown(); + environment = null; + } + } + } + + long getSampled() { + return sampled; + } + + /** + * Counts this fixture's live checkpoint scheduler threads. This is the cost a shared scheduler + * exists to remove, so it is reported rather than derived. + */ + public long countSchedulerThreads() { + return schedulerThreads().size(); + } + + /** + * Live checkpoint scheduler threads started since this fixture's setup began. Threads that + * already existed, such as ones another test in the same JVM has not finished stopping, are not + * counted. + */ + List<Thread> schedulerThreads() { + return liveSchedulerThreads().stream() + .filter(thread -> !preexistingSchedulerThreads.contains(thread)) + .collect(Collectors.toList()); + } + + private static List<Thread> liveSchedulerThreads() { + return Thread.getAllStackTraces().keySet().stream() + .filter(thread -> thread.getName().startsWith(SCHEDULER_THREAD_NAME_PREFIX)) + .collect(Collectors.toList()); + } + + private void createManagers(SeaTunnelServer server, CheckpointConfig memberConfig) { + NodeEngine nodeEngine = server.getNodeEngine(); + IMap<Object, Object> runningJobState = + nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_RUNNING_JOB_STATE); + CheckpointConfig jobConfig = new CheckpointConfig(); + jobConfig.setCheckpointInterval(intervalMillis); + jobConfig.setStorage(memberConfig.getStorage()); + + pendingCounters = new AtomicInteger[pipelineNum]; + lastTriggerNanos = new long[pipelineNum]; + for (int i = 0; i < pipelineNum; i++) { + long jobId = FIRST_JOB_ID + i; + FakeTaskCheckpointManager manager = + new FakeTaskCheckpointManager( + jobId, + nodeEngine, + singleTaskPlan(jobId), + jobConfig, + server.getCheckpointService().getCheckpointStorage(), + coordinatorExecutor, + runningJobState, + server.getEngineContext(), + server.getCheckpointMonitorService(), + failure); + managers.add(manager); + pendingCounters[i] = readPendingCounter(manager.getCheckpointCoordinator(PIPELINE_ID)); + lastTriggerNanos[i] = NOT_SYNCED; + } + } + + private void startStaggered() { + long start = System.nanoTime(); + for (int i = 0; i < pipelineNum; i++) { + parkUntil(start + intervalNanos * i / pipelineNum); + managers.get(i).startTask(); + } + } + + /** + * Spins over every {@code stride}-th coordinator until each one without a known phase has been + * seen idle and then triggering. Those that already have a phase are refreshed whenever they + * trigger during the wait, so they do not fall out of the rotation meanwhile. Not measured; it + * occupies one core for up to a few intervals. + * + * @param stride {@code probeStride} to cover the probes, 1 to cover every coordinator + */ + private void resync(int stride) { + boolean[] seenIdle = new boolean[pipelineNum]; + long deadline = System.nanoTime() + RESYNC_INTERVALS * intervalNanos + PENDING_REARM_NANOS; + long covered = strideCount(stride); + long unsynced = covered - syncedCount(stride); + while (unsynced > 0) { + checkFailure(); + if (System.nanoTime() > deadline) { + throw new IllegalStateException( + unsynced + + " of " + + covered + + " checkpoint coordinators did not trigger within " + + RESYNC_INTERVALS + + " intervals plus one pending re-arm"); + } + for (int i = 0; i < pipelineNum; i += stride) { + if (pendingCounters[i].get() == 0) { + seenIdle[i] = true; + } else if (seenIdle[i]) { + if (lastTriggerNanos[i] == NOT_SYNCED) { + unsynced--; + } + lastTriggerNanos[i] = System.nanoTime(); + seenIdle[i] = false; + } + } + } + } + + /** + * Returns the probe with a known phase that is due soonest. Resynchronises first once fewer + * than half the probes have a known phase: skipped probes would otherwise stay out of the + * rotation for the rest of the iteration, and with few probes a single skip can leave none. + */ + private int soonestSynced() { + if (syncedCount(probeStride) * 2 < probeCount()) { + resync(probeStride); + } + int soonest = -1; + for (int i = 0; i < pipelineNum; i += probeStride) { + if (lastTriggerNanos[i] != NOT_SYNCED + && (soonest < 0 || lastTriggerNanos[i] < lastTriggerNanos[soonest])) { + soonest = i; + } + } + return soonest; + } + + public long probeCount() { + return strideCount(probeStride); + } + + private long strideCount(int stride) { + return (pipelineNum + stride - 1) / stride; + } + + private long syncedCount(int stride) { + long synced = 0; + for (int i = 0; i < pipelineNum; i += stride) { + if (lastTriggerNanos[i] != NOT_SYNCED) { + synced++; + } + } + return synced; + } + + private void widenSpinWindow(long parkOvershootNanos) { + spinWindowNanos = + Math.min( + MAX_SPIN_WINDOW_NANOS, + Math.max(spinWindowNanos, parkOvershootNanos + MIN_SPIN_WINDOW_NANOS)); + } + + private void skipPending(int coordinator) { + lastTriggerNanos[coordinator] = NOT_SYNCED; + skippedPending++; + } + + private void skipCollided(int coordinator) { + lastTriggerNanos[coordinator] = NOT_SYNCED; + skippedCollided++; + } + + private void skipOverrun(int coordinator) { + lastTriggerNanos[coordinator] = NOT_SYNCED; + skippedOverrun++; + } + + private void skipEarly(int coordinator) { + lastTriggerNanos[coordinator] = NOT_SYNCED; + skippedEarly++; + } + + private void checkFailure() { + Throwable throwable = failure.get(); + if (throwable != null) { + throw new IllegalStateException("Checkpoint coordinator failed", throwable); + } + } + + private void validateParameters() { + if (pipelineNum < 1) { + throw new IllegalArgumentException("pipelineNum must be at least 1"); + } + if (intervalMillis < MIN_CHECKPOINT_INTERVAL_MILLIS) { + throw new IllegalArgumentException( + "checkpoint interval must be at least " + + MIN_CHECKPOINT_INTERVAL_MILLIS + + " ms to match the minimum SeaTunnel accepts"); + } + } + + /** + * Starts are spread evenly across one interval, so coordinators {@code i} and {@code i + + * stride} start {@code stride * interval / pipelineNum} apart. Returns the smallest stride that + * keeps that at least {@link #MIN_PROBE_SPACING_NANOS}, and at most {@code pipelineNum} so + * there is always one probe. + */ + static int probeStride(int pipelineNum, long intervalNanos) { + long stride = (MIN_PROBE_SPACING_NANOS * pipelineNum + intervalNanos - 1) / intervalNanos; + return (int) Math.max(1L, Math.min(stride, pipelineNum)); + } + + private static CheckpointPlan singleTaskPlan(long jobId) { + TaskLocation task = new TaskLocation(new TaskGroupLocation(jobId, PIPELINE_ID, 1L), 0, 0); + return CheckpointPlan.builder() + .pipelineId(PIPELINE_ID) + .pipelineSubtasks(Collections.singleton(task)) + .startingSubtasks(Collections.singleton(task)) + .pipelineActions(Collections.emptyMap()) + .subtaskActions(Collections.emptyMap()) + .build(); + } + + /** Same shape as the executor {@code CoordinatorService} hands every {@code JobMaster}. */ + private static ThreadPoolExecutor createCoordinatorExecutor(EngineConfig engineConfig) { + AtomicInteger threadIndex = new AtomicInteger(); + return new ThreadPoolExecutor( + engineConfig.getCoordinatorServiceConfig().getCoreThreadNum(), + engineConfig.getCoordinatorServiceConfig().getMaxThreadNum(), + EXECUTOR_KEEP_ALIVE_SECONDS, + TimeUnit.SECONDS, + new SynchronousQueue<>(), + runnable -> { + Thread thread = new Thread(runnable); + thread.setName( + "benchmark-coordinator-service-" + threadIndex.getAndIncrement()); + return thread; + }); + } + + private static AtomicInteger readPendingCounter(CheckpointCoordinator coordinator) { + try { + return (AtomicInteger) PENDING_COUNTER_FIELD.get(coordinator); + } catch (IllegalAccessException e) { + throw new IllegalStateException("Cannot read the pending checkpoint counter", e); + } + } + + /** + * The storage benchmarks' member, with its checkpoint storage isolated in the trial's temporary + * directory, but without IMap persistence. Its write-through MapStore puts a file write on a + * partition thread for every engine IMap update, including checkpoint id allocation, which + * stalls checkpoints for hundreds of milliseconds and is not what this benchmark measures. + */ + private static final class SchedulingEnvironmentContext + extends SeaTunnelStorageEnvironmentContext { + + private static final String ENGINE_MAPS = "engine*"; + + @Override + protected SeaTunnelConfig createSeaTunnelConfig(String clusterName) { + SeaTunnelConfig config = super.createSeaTunnelConfig(clusterName); + config.getHazelcastConfig() + .getMapConfig(ENGINE_MAPS) + .getMapStoreConfig() + .setEnabled(false); + return config; + } + } + + private static void parkUntil(long deadlineNanos) { + long remaining; + while ((remaining = deadlineNanos - System.nanoTime()) > 0) { + LockSupport.parkNanos(remaining); + } + } +} diff --git a/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/FakeTaskCheckpointManager.java b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/FakeTaskCheckpointManager.java new file mode 100644 index 0000000000..d57bed4a45 --- /dev/null +++ b/seatunnel-benchmarks/src/main/java/org/apache/seatunnel/benchmark/checkpoint/FakeTaskCheckpointManager.java @@ -0,0 +1,160 @@ +/* + * 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.seatunnel.benchmark.checkpoint; + +import org.apache.seatunnel.engine.checkpoint.storage.api.CheckpointStorage; +import org.apache.seatunnel.engine.common.config.server.CheckpointConfig; +import org.apache.seatunnel.engine.server.checkpoint.CheckpointBarrier; +import org.apache.seatunnel.engine.server.checkpoint.CheckpointManager; +import org.apache.seatunnel.engine.server.checkpoint.CheckpointPlan; +import org.apache.seatunnel.engine.server.checkpoint.monitor.CheckpointMonitorService; +import org.apache.seatunnel.engine.server.checkpoint.operation.CheckpointBarrierTriggerOperation; +import org.apache.seatunnel.engine.server.checkpoint.operation.TaskAcknowledgeOperation; +import org.apache.seatunnel.engine.server.checkpoint.operation.TaskReportStatusOperation; +import org.apache.seatunnel.engine.server.common.SeaTunnelEngineContext; +import org.apache.seatunnel.engine.server.execution.TaskLocation; +import org.apache.seatunnel.engine.server.task.operation.TaskOperation; +import org.apache.seatunnel.engine.server.task.statemachine.SeaTunnelTaskState; +import org.apache.seatunnel.engine.server.utils.NodeEngineUtil; + +import com.hazelcast.map.IMap; +import com.hazelcast.spi.impl.NodeEngine; +import com.hazelcast.spi.impl.operationservice.Operation; +import com.hazelcast.spi.impl.operationservice.impl.InvocationFuture; + +import java.lang.reflect.Field; +import java.util.Collections; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.atomic.AtomicReference; + +/** + * A real {@link CheckpointManager} whose tasks are fakes answered at the network boundary. + * + * <p>Everything on the coordinator side is production code: the coordinators, their scheduling, + * checkpoint id allocation, the IMap state and the checkpoint storage. Only the messages that would + * travel to a task are intercepted, and the fake answers exactly two of them: + * + * <ul> + * <li>{@link #startTask}: the task reports {@code READY_START}, which is what makes the + * coordinator arm its periodic trigger. + * <li>{@link #sendOperationToMemberNode} with a {@link CheckpointBarrierTriggerOperation}: every + * task of the pipeline acknowledges the barrier at once, with no state. + * </ul> + * + * <p>Every other operation is answered with an empty success. If the task protocol changes, for + * example a new message the coordinator waits on before arming the trigger or before completing a + * checkpoint, this class is what needs updating, and the fixture then fails at setup rather than + * reporting numbers from a coordinator that never triggers. + */ +final class FakeTaskCheckpointManager extends CheckpointManager { + + /** + * {@code CheckpointBarrierTriggerOperation} has no getter for its barrier, and the barrier is + * needed to acknowledge it. Read-only; if the field is renamed or removed, loading this class + * fails with a message naming it instead of the benchmark silently never completing a + * checkpoint. + */ + private static final Field BARRIER_FIELD = + BenchmarkReflection.requireField(CheckpointBarrierTriggerOperation.class, "barrier"); + + private final long jobId; + private final NodeEngine nodeEngine; + private final CheckpointPlan plan; + private final AtomicReference<Throwable> failure; + + FakeTaskCheckpointManager( + long jobId, + NodeEngine nodeEngine, + CheckpointPlan plan, + CheckpointConfig checkpointConfig, + CheckpointStorage checkpointStorage, + ExecutorService executorService, + IMap<Object, Object> runningJobStateIMap, + SeaTunnelEngineContext engineContext, + CheckpointMonitorService checkpointMonitorService, + AtomicReference<Throwable> failure) { + super( + jobId, + false, + null, + null, + nodeEngine, + null, + Collections.singletonMap(plan.getPipelineId(), plan), + checkpointConfig, + checkpointStorage, + executorService, + runningJobStateIMap, + engineContext, + checkpointMonitorService); + this.jobId = jobId; + this.nodeEngine = nodeEngine; + this.plan = plan; + this.failure = failure; + } + + /** + * Starts the pipeline the way a deployed job does: the pipeline is reported running, then every + * task reports {@code READY_START}. The coordinator arms its periodic trigger one checkpoint + * interval after the last report. + */ + void startTask() { + reportedPipelineRunning(plan.getPipelineId(), false); + for (TaskLocation task : plan.getPipelineSubtasks()) { + reportedTask(new TaskReportStatusOperation(task, SeaTunnelTaskState.READY_START)); + } + } + + @Override + protected InvocationFuture<?> sendOperationToMemberNode(TaskOperation operation) { + if (operation instanceof CheckpointBarrierTriggerOperation) { + acknowledgeBarrier(readBarrier((CheckpointBarrierTriggerOperation) operation)); + } + return NodeEngineUtil.sendOperationToMemberNode( + nodeEngine, new NoOpOperation(), nodeEngine.getThisAddress()); + } + + /** Records the failure instead of reaching the absent {@code JobMaster}. */ + @Override + protected void handleCheckpointError(int pipelineId, boolean neverRestore) { + failure.compareAndSet( + null, + new IllegalStateException( + "Checkpoint coordinator of job " + jobId + " reported an error")); + } + + private void acknowledgeBarrier(CheckpointBarrier barrier) { + for (TaskLocation task : plan.getPipelineSubtasks()) { + acknowledgeTask(new TaskAcknowledgeOperation(task, barrier, Collections.emptyList())); + } + } + + private static CheckpointBarrier readBarrier(CheckpointBarrierTriggerOperation operation) { + try { + return (CheckpointBarrier) BARRIER_FIELD.get(operation); + } catch (IllegalAccessException e) { + throw new IllegalStateException("Cannot read the checkpoint barrier", e); + } + } + + /** Completes on the local member without doing anything, standing in for a task's reply. */ + private static final class NoOpOperation extends Operation { + @Override + public void run() {} + } +} diff --git a/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/CheckpointSchedulingBenchmarkTest.java b/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/CheckpointSchedulingBenchmarkTest.java new file mode 100644 index 0000000000..a9217bde95 --- /dev/null +++ b/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/CheckpointSchedulingBenchmarkTest.java @@ -0,0 +1,42 @@ +/* + * 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.seatunnel.benchmark; + +import org.junit.jupiter.api.Test; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Threads; + +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class CheckpointSchedulingBenchmarkTest { + + @Test + void shouldSampleSchedulingDelayOnASingleThread() { + assertEquals( + Mode.SampleTime, + CheckpointSchedulingBenchmark.class.getAnnotation(BenchmarkMode.class).value()[0]); + assertEquals( + TimeUnit.MICROSECONDS, + CheckpointSchedulingBenchmark.class.getAnnotation(OutputTimeUnit.class).value()); + assertEquals(1, CheckpointSchedulingBenchmark.class.getAnnotation(Threads.class).value()); + } +} diff --git a/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/checkpoint/CheckpointSchedulingFixtureTest.java b/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/checkpoint/CheckpointSchedulingFixtureTest.java new file mode 100644 index 0000000000..c28b4786b2 --- /dev/null +++ b/seatunnel-benchmarks/src/test/java/org/apache/seatunnel/benchmark/checkpoint/CheckpointSchedulingFixtureTest.java @@ -0,0 +1,150 @@ +/* + * 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.seatunnel.benchmark.checkpoint; + +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class CheckpointSchedulingFixtureTest { + + private static final int PIPELINE_NUM = 4; + private static final long CHECKPOINT_INTERVAL_MILLIS = 200L; + private static final int SAMPLE_COUNT = 20; + private static final long THREAD_STOP_TIMEOUT_NANOS = TimeUnit.SECONDS.toNanos(30); + + @Test + void shouldMeasureDueTriggersOfRealCoordinators() throws Exception { + CheckpointSchedulingFixture fixture = + new CheckpointSchedulingFixture(PIPELINE_NUM, CHECKPOINT_INTERVAL_MILLIS); + fixture.setUp(); + try { + assertTrue( + fixture.countSchedulerThreads() > 0, + "the coordinators should be running checkpoint scheduler threads"); + fixture.beginIteration(); + for (int i = 0; i < SAMPLE_COUNT; i++) { + fixture.awaitNextDueTrigger(); + long due = System.nanoTime(); + long observed = fixture.awaitTrigger(); + assertTrue( + observed - due < TimeUnit.MILLISECONDS.toNanos(CHECKPOINT_INTERVAL_MILLIS), + "a due trigger should run well within one interval"); + } + + assertEquals(SAMPLE_COUNT, fixture.getSampled()); + fixture.endIteration(); + } finally { + fixture.tearDown(); + } + } + + @Test + void shouldStopEverySchedulerThreadOnTearDown() throws Exception { + // Stands in for a scheduler thread another test in this JVM has not finished stopping. + CountDownLatch release = new CountDownLatch(1); + Thread foreign = new Thread(() -> awaitQuietly(release), "checkpoint-foreign"); + foreign.setDaemon(true); + foreign.start(); + try { + CheckpointSchedulingFixture fixture = + new CheckpointSchedulingFixture(PIPELINE_NUM, CHECKPOINT_INTERVAL_MILLIS); + fixture.setUp(); + fixture.tearDown(); + + assertSchedulerThreadsStop(fixture); + assertTrue(foreign.isAlive(), "the foreign thread should not have been waited for"); + } finally { + release.countDown(); + } + } + + private static void assertSchedulerThreadsStop(CheckpointSchedulingFixture fixture) + throws InterruptedException { + long deadline = System.nanoTime() + THREAD_STOP_TIMEOUT_NANOS; + while (fixture.countSchedulerThreads() > 0 && System.nanoTime() < deadline) { + TimeUnit.MILLISECONDS.sleep(CHECKPOINT_INTERVAL_MILLIS); + } + assertEquals( + 0L, + fixture.countSchedulerThreads(), + () -> + "still running: " + + fixture.schedulerThreads().stream() + .map(Thread::getName) + .collect(Collectors.toList())); + } + + private static void awaitQuietly(CountDownLatch latch) { + try { + latch.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + @Test + void shouldRejectParametersOutsideTheSupportedRange() { + assertThrows( + IllegalArgumentException.class, + () -> new CheckpointSchedulingFixture(0, 1_000L).setUp()); + assertThrows( + IllegalArgumentException.class, + () -> new CheckpointSchedulingFixture(1, 9L).setUp()); + } + + @Test + void shouldSpaceMeasuredCoordinatorsAtLeastOneHundredMillisApart() { + long interval = TimeUnit.MILLISECONDS.toNanos(CHECKPOINT_INTERVAL_MILLIS); + + assertEquals(1, CheckpointSchedulingFixture.probeStride(1, interval)); + assertEquals(2, CheckpointSchedulingFixture.probeStride(4, interval)); + assertEquals(5, CheckpointSchedulingFixture.probeStride(10, interval)); + assertEquals(5, CheckpointSchedulingFixture.probeStride(500, TimeUnit.SECONDS.toNanos(10))); + assertEquals( + 3, CheckpointSchedulingFixture.probeStride(3, TimeUnit.MILLISECONDS.toNanos(10))); + } + + @Test + void shouldRejectTooManySkipsOnlyOnceEnoughTriggersWereDue() { + // A one-second smoke iteration sees a handful of due triggers; one skip there is noise. + assertFalse(CheckpointSchedulingFixture.isSkipShareTooHigh(7, 1)); + assertFalse(CheckpointSchedulingFixture.isSkipShareTooHigh(49, 49)); + + assertFalse(CheckpointSchedulingFixture.isSkipShareTooHigh(100, 10)); + assertTrue(CheckpointSchedulingFixture.isSkipShareTooHigh(100, 11)); + assertTrue(CheckpointSchedulingFixture.isSkipShareTooHigh(50, 50)); + } + + @Test + void shouldNameTheEngineFieldWhenItNoLongerExists() { + IllegalStateException failure = + assertThrows( + IllegalStateException.class, + () -> BenchmarkReflection.requireField(Object.class, "pendingCounter")); + + assertTrue(failure.getMessage().contains("java.lang.Object#pendingCounter")); + } +} diff --git a/tools/benchmarks/save_jmh_result.py b/tools/benchmarks/save_jmh_result.py index acbca278c8..fafa3a63f2 100644 --- a/tools/benchmarks/save_jmh_result.py +++ b/tools/benchmarks/save_jmh_result.py @@ -78,6 +78,25 @@ def flatten(values): return [value for fork in values for value in fork] +def histogram_mean(buckets): + count = sum(bucket_count for _, bucket_count in buckets) + return sum(value * bucket_count for value, bucket_count in buckets) / count + + +def iteration_scores(primary): + # Sample modes publish rawDataHistogram (fork -> iteration -> [value, count]) instead of + # rawData. Each iteration's mean keeps the samples comparable with the other modes: they + # describe variation between iterations, not the spread of individual invocations. + if "rawData" in primary: + return flatten(primary["rawData"]) + return [ + histogram_mean(iteration) + for fork in primary.get("rawDataHistogram", []) + for iteration in fork + if iteration + ] + + def benchmark_name(result): params = result.get("params", {}) suffix = ",".join("{}={}".format(key, params[key]) for key in sorted(params)) @@ -90,7 +109,7 @@ def jmh_metrics(results): primary = result["primaryMetric"] score = finite_or_none(primary.get("score")) error = finite_or_none(primary.get("scoreError")) - samples = flatten(primary.get("rawData", [])) + samples = iteration_scores(primary) metrics.append( { "name": benchmark_name(result), diff --git a/tools/benchmarks/test_save_jmh_result.py b/tools/benchmarks/test_save_jmh_result.py index fa107a5372..16c2ec1f0b 100644 --- a/tools/benchmarks/test_save_jmh_result.py +++ b/tools/benchmarks/test_save_jmh_result.py @@ -57,6 +57,32 @@ class SaveJmhResultTest(unittest.TestCase): self.assertEqual(0.05, metric["relative_score_error"]) self.assertEqual("higher", metric["direction"]) + def test_uses_iteration_means_of_sample_time_histograms(self): + metrics = save_jmh_result.jmh_metrics( + [ + { + "benchmark": "org.apache.seatunnel.Checkpoint.triggerDelay", + "mode": "sample", + "forks": 2, + "params": {}, + "primaryMetric": { + "score": 20.0, + "scoreError": 2.0, + "scoreUnit": "us/op", + "rawDataHistogram": [ + [[[10.0, 3], [30.0, 1]], [[20.0, 2]]], + [[[25.0, 4]], []], + ], + }, + } + ] + ) + + metric = metrics[0] + self.assertEqual([15.0, 20.0, 25.0], metric["samples"]) + self.assertEqual(5.0, metric["sample_standard_deviation"]) + self.assertEqual("lower", metric["direction"]) + def test_aggregates_pipeline_medians_correctness_and_clamping(self): with tempfile.TemporaryDirectory() as directory: pipeline_dir = pathlib.Path(directory)
