nzw921rx opened a new pull request, #12044:
URL: https://github.com/apache/seatunnel/pull/12044
### Purpose of this pull request
This PR adds JMH coverage for checkpoint storage and the production
file-backed IMap storage hot paths used by Zeta. The benchmarks measure real
storage behavior under configurable state pressure instead of timing isolated
API calls or mock implementations.
#### Benchmark environment and fixture design
- `SeaTunnelStorageEnvironmentContext` starts an isolated embedded
Zeta/Hazelcast environment with IMap MapStore enabled, backup count set to `0`,
and local filesystem-backed storage. It owns only the runtime lifecycle and
exposes the production storage services required by workloads.
- `StorageLifecycleFixtureJob` submits a real streaming job and captures
production `JobInfo`, job history, metrics, and completed-checkpoint values.
Job startup, checkpoint coordination/barrier/ACK time, and fixture generation
stay outside the measured storage operations.
- `BenchmarkCheckpointStorageFactory` decorates the production `HdfsStorage`
only while generating checkpoint fixtures. The measured checkpoint writes use a
real `HdfsStorage` configured with `file:///`.
- `JobDagFixtureFactory` builds production `JobDAGInfo` values in code with
an exact pipeline count, so DAG payload size can be controlled without parsing
a fixture job configuration.
- `BenchmarkTemplates` centralizes UTF-8 classpath template loading and
placeholder rendering. Several benchmark environments need the same behavior;
extracting it removes context-specific wrappers, applies consistent
missing-resource/unresolved-placeholder validation, and keeps environment
contexts focused on lifecycle management.
#### Added benchmark scenarios
| Benchmark | Scenario | Production hot path | What it measures |
|---|---|---|---|
| `CheckpointStorageBenchmark` | `checkpointPersistenceTransaction` |
`StateStoreCheckpointIDCounter.getAndIncrement` →
`CheckpointMonitorService.onCheckpointCompleted` →
`HdfsStorage.storeCheckPoint` | End-to-end storage portion of one completed
checkpoint: atomic ID allocation, overview update, serialization,
temporary-file write, and atomic rename. Checkpoint barrier/ACK time is
excluded. |
| `CheckpointStorageBenchmark` | `checkpointIdAtomicIncrement` | checkpoint
counter state store | Atomic checkpoint ID allocation cost in isolation. |
| `CheckpointStorageBenchmark` | `checkpointOverviewIncrementalUpdate` |
checkpoint overview IMap compute/update | Incremental checkpoint metadata
update using a coordinator-produced `CompletedCheckpoint`. |
| `IMapJobStorageBenchmark` | `taskGroupStateTransition` | running-state and
state-timestamp IMaps used by task state transitions | Repeated timestamp
get/set and task-group state get/set/get under configurable stored-task
pressure. |
| `IMapJobStorageBenchmark` | `runningMetricsReport` |
`ReportMetricsOperation` → metrics snapshot merge | The periodic
worker-to-master metrics reporting path used by `TaskExecutionService`,
including Hazelcast operation dispatch and IMap merge/compute work. |
| `IMapJobStorageBenchmark` | `runningJobGrowth` | running `JobInfo`, task
state, and timestamp IMaps | Cost as retained running-job cardinality
continuously grows across invocations. |
| `IMapJobStorageBenchmark` | `completedJobHistoryGrowth` |
`JobHistoryService.storeFinishedPipelineMetrics` / `storeFinishedJobState` plus
running-state cleanup | Normal completed-job persistence while finished state
and metrics accumulate until production TTL expiration. |
| `IMapJobStorageBenchmark` | `runningJobRecovery` | `IMap.loadAll(true)` →
`FileMapStore` / WAL reader → `entrySet` | Full running-job metadata reload and
scan after in-memory entries are evicted, matching master recovery behavior. |
| `IMapDagStorageBenchmark` | `finishedJobDagStore` | finished-job DAG IMap
→ synchronous `FileMapStore` WAL append | DAG persistence cost across 1/10/100
code-built pipelines and configurable existing DAG pressure. |
| `IMapDagStorageBenchmark` | `finishedJobDagLoad` | `IMap.loadAll(keys,
true)` → `FileMapStore` / WAL reader | Reloading an evicted production
`JobDAGInfo`. |
| `IMapWalStorageBenchmark` | `appendNewKey` | finished-job DAG `IMap.put` →
`FileMapStore` → WAL writer | WAL append cost and actual WAL-byte growth when
key cardinality grows. |
| `IMapWalStorageBenchmark` | `appendHotKey` | repeated `IMap.put` for one
key → WAL writer | Hot-key mutation cost and WAL-byte growth for update-heavy
histories. |
| `IMapWalStorageBenchmark` | `recoverAll` | `IMap.loadAll(true)` →
`FileMapStore` → WAL reader | Recovery scaling when live key count and
mutations per key vary independently, including the deep-history pattern
targeted by #11494. WAL construction is outside measured time. |
All measured paths use production Hazelcast IMaps, `FileMapStore`, WAL
serialization/reader/writer, checkpoint state stores, and HDFS checkpoint
storage over `file:///`; no mock storage is used.
#### Scheduled core suite versus manual execution
The workflow has a 240-minute limit and runs both Java 8 and Java 11.
Running every parameter combination with the standard 3 forks, 3 warmup
iterations, and 5 measurement iterations would make scheduled runs too large
and less reliable. Therefore:
- scheduled runs set `BENCHMARK_SUITE=benchmarks_core` and resolve the exact
stable hot-path selectors from `tools/benchmarks/suites/benchmarks_core.txt`;
- manual runs keep class-level choices for broader investigation, while
`custom_benchmarks` remains available for any exact JMH regular expression;
- `run_benchmarks.sh` validates the suite name/file, converts non-comment
suite entries into the JMH selector, and records the suite in generated reports;
- method-level workflow choices such as `sourceSink$` were removed because
every additional benchmark method would grow and duplicate the dropdown.
`sourceSink` and `sourceTransformSink` are still explicitly pinned in the core
suite, while users can select the whole `SeaTunnelPipelineBenchmark` class or
enter an exact method through `custom_benchmarks`.
### Does this PR introduce _any_ user-facing change?
No. This PR only adds benchmark infrastructure, benchmark scenarios, tests,
and benchmark workflow selection.
### How was this patch tested?
Formatting:
```bash
./mvnw -f seatunnel-benchmarks/pom.xml spotless:apply
```
Benchmark module and dependency build:
```bash
./mvnw -Pbenchmark -pl seatunnel-benchmarks -am -DskipTests package
./mvnw -f seatunnel-benchmarks/pom.xml package -DskipTests
```
New benchmark, fixture, and template tests (14 tests, 0 failures):
```bash
./mvnw -Pbenchmark -pl seatunnel-benchmarks \
-Dtest=BenchmarkTemplatesTest,JobDagFixtureFactoryTest,CheckpointStorageBenchmarkTest,IMapJobStorageBenchmarkTest,IMapDagStorageBenchmarkTest,IMapWalStorageBenchmarkTest
\
test
```
Benchmark report/script tests (23 tests, 0 failures):
```bash
python3 -B -m unittest discover -s tools/benchmarks -p 'test_*.py'
```
Workflow and script syntax:
```bash
bash -n tools/benchmarks/run_benchmarks.sh
ruby -e 'require "yaml"; YAML.load_file(".github/workflows/benchmarks.yml")'
```
The core-suite selectors were also resolved against the built JMH jar, and a
forked `CheckpointingTimeBenchmark.checkpointSingleInput` smoke run completed
successfully against the real checkpoint environment.
### Check list
* [x] No new Jar binary package is added.
* [x] No user-facing documentation change is required for benchmark-only
code.
* [x] No incompatible change is introduced.
* [x] This PR does not modify connector code.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]