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)

Reply via email to