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]

Reply via email to