nzw921rx commented on issue #11598:
URL: https://github.com/apache/seatunnel/issues/11598#issuecomment-5507879596
## Architecture Overview
This section summarizes the final architecture of the Zeta benchmark
framework after #11671, #11986, #12021, #12026, and #12044.
The goal is not to benchmark every Engine class.
The framework provides a reusable way to:
> measure a real performance path → compare revisions → diagnose a change →
verify an optimization.
---
### 1. Overall architecture
The benchmark framework lives in the main SeaTunnel repository as an opt-in
`seatunnel-benchmarks` module.
It directly exercises production Zeta implementations while remaining
outside the normal build and runtime classpath.
```mermaid
flowchart TB
CI["Execution<br/>GitHub Actions / Local"]
JMH["JMH benchmark layer"]
Env["Benchmark environments<br/>& deterministic fixtures"]
Zeta["Production Zeta code"]
Result["Results & comparison"]
Diag["Diagnostics"]
CI --> JMH
JMH --> Env
Env --> Zeta
JMH --> Result
Zeta --> Result
Result --> Diag
```
The main boundary is:
```text
benchmark-only infrastructure
↓
production implementation
```
Benchmark convenience should not leak into production APIs.
When internal access is required, prefer a narrowly scoped benchmark-side
bridge rather than adding an Engine API only for testing.
---
### 2. Benchmark coverage
The current suite covers the major performance-sensitive Zeta paths.
```mermaid
flowchart LR
Record["Record processing"]
Row["SeaTunnelRow"]
Queue["Intermediate Queue"]
Pipeline["Complete Pipeline"]
Checkpoint["Checkpoint"]
Coord["Coordination"]
CPStore["Persistence"]
State["Engine State"]
IMap["IMap / StateStore"]
DAG["Job DAG"]
WAL["WAL / Recovery"]
Record --> Row
Record --> Queue
Record --> Pipeline
Checkpoint --> Coord
Checkpoint --> CPStore
State --> IMap
State --> DAG
State --> WAL
```
Current implementations:
| Area | Benchmark |
| --- | --- |
| Row operations | `SeaTunnelRowBenchmark` |
| Record handoff | `IntermediateQueueBenchmark` |
| End-to-end data path | `SeaTunnelPipelineBenchmark` |
| Checkpoint coordination | `CheckpointingTimeBenchmark` |
| Checkpoint persistence | `CheckpointStorageBenchmark` |
| Job/metrics state | `IMapJobStorageBenchmark` |
| DAG persistence | `IMapDagStorageBenchmark` |
| WAL and recovery | `IMapWalStorageBenchmark` |
New benchmark areas should be added only when an actual performance question
requires them.
---
### 3. Measurement model
A common rule across the framework is:
> Setup and validation should not accidentally become part of the measured
operation.
```mermaid
flowchart LR
Setup["Prepare environment<br/>and fixtures"]
Measure["Measured production path"]
Verify["Validate correctness<br/>and persistence"]
Cleanup["Cleanup"]
Setup --> Measure --> Verify --> Cleanup
style Measure stroke-width:4px
```
The exact measurement boundary depends on the question.
For example:
```mermaid
flowchart TB
Pipeline["Pipeline benchmark<br/>Job submission → complete data path →
job completion"]
CP["Checkpoint benchmark<br/>Trigger → barrier/snapshot/ACK →
persistence → completion"]
Storage["Storage benchmark<br/>Prepared production fixture → target
persistence operation"]
Queue["Queue benchmark<br/>One production record handoff"]
```
This prevents unrelated startup or fixture cost from hiding the path being
investigated.
---
### 4. Pipeline workload model
The full Pipeline benchmark uses deterministic in-memory Source, Transform,
and Sink implementations over a real embedded Zeta runtime.
```mermaid
flowchart LR
Source["BenchmarkSource<br/>open-loop schedule"]
Queue1["Zeta runtime"]
Transform["Transform<br/>optional"]
Queue2["Zeta runtime"]
Sink["BenchmarkSink"]
Metrics["Throughput<br/>P50/P95/P99<br/>latency growth"]
Source --> Queue1 --> Transform --> Queue2 --> Sink
Sink --> Metrics
```
The Source follows an absolute open-loop schedule.
When Zeta cannot keep up, the Source schedule continues advancing, so
backlog appears as increasing end-to-end latency instead of being hidden by a
slower producer.
Correctness is checked before interpreting performance:
```text
expected rows == processed rows
Transform checksum is valid
```
---
### 5. Persistent-state benchmark model
Storage benchmarks use real Zeta StateStore, Hazelcast IMap, `FileMapStore`,
WAL, and checkpoint storage with local files.
Persistent mutation benchmarks use **fixed logical work**, rather than “as
many operations as possible”.
```mermaid
flowchart LR
Fixture["Prepare identical state"]
A["Revision A<br/>100 operations"]
B["Revision B<br/>100 operations"]
Compare["Compare time/op"]
Fixture --> A --> Compare
Fixture --> B --> Compare
```
This is important because a throughput-style storage benchmark could
otherwise make a faster revision create more:
```text
keys
WAL records
files
mutation history
```
and therefore change its own workload.
Fixed batches plus `OperationsPerInvocation` keep storage pressure
comparable.
---
### 6. Stable core vs investigation benchmarks
Not every benchmark belongs in scheduled execution.
```mermaid
flowchart TB
All["All benchmark scenarios"]
Core["Stable core suite<br/>scheduled"]
Manual["Focused benchmarks<br/>manual"]
Expensive["Expensive diagnostics<br/>manual only"]
All --> Core
All --> Manual
All --> Expensive
```
`benchmarks_core` contains representative, stable, bounded workloads.
Large parameter matrices such as deep WAL recovery remain on-demand.
The core suite should stay small enough that scheduled runs remain practical.
---
### 7. Execution and comparison
Normal execution produces raw and normalized artifacts rather than only
console output.
```mermaid
flowchart LR
Run["JMH run"]
Raw["Raw JMH JSON"]
Pipe["Pipeline JSON"]
Normalize["Normalized report"]
Summary["Markdown summary"]
Run --> Raw
Run --> Pipe
Raw --> Normalize
Pipe --> Normalize
Normalize --> Summary
```
The normalized format records:
```text
benchmark
parameters
score / direction
raw samples
commit
JDK / JVM
runner
CPU / memory
workload identity
correctness
```
This keeps the current framework usable even if historical storage or
visualization is added later.
---
### 8. PR comparison
PR comparison uses the same worker and alternates revisions:
```mermaid
flowchart LR
B1["Baseline"]
C1["Candidate"]
C2["Candidate"]
B2["Baseline"]
B1 --> C1 --> C2 --> B2
```
The report compares the median of the two baseline runs with the median of
the two candidate runs.
Metric direction is normalized so:
> positive change means improvement
for both higher-is-better throughput and lower-is-better latency/time
metrics.
The result remains observational because GitHub-hosted runner performance is
not sufficiently controlled for a hard regression gate.
---
### 9. Measurement and diagnostics are separate
Profiling is used to explain a benchmark result, not to produce the
benchmark baseline itself.
```mermaid
flowchart LR
Benchmark["Normal benchmark"]
Change["Unexpected difference"]
CPU["CPU"]
Wall["Wall"]
Lock["Lock"]
GC["GC / Allocation"]
JFR["JFR"]
Root["Root cause"]
Fix["Production optimization"]
Verify["Benchmark again"]
Benchmark --> Change
Change --> CPU
Change --> Wall
Change --> Lock
Change --> GC
Change --> JFR
CPU --> Root
Wall --> Root
Lock --> Root
GC --> Root
JFR --> Root
Root --> Fix --> Verify
```
Diagnostics are intentionally kept in a separate workflow because profiler
overhead changes benchmark behavior.
Supported diagnostics currently include:
```text
CPU
Wall clock
Lock contention
GC / allocation
JFR
```
---
### 10. How the framework should evolve
The framework is now intended to grow **from real performance questions**,
rather than from a predefined list of classes.
```mermaid
flowchart LR
Change["Engine change<br/>or suspicious result"]
Existing{"Existing benchmark<br/>can measure it?"}
Use["Use existing benchmark"]
Add["Add focused benchmark"]
Diagnose["Compare / diagnose"]
Optimize["Optimize Zeta"]
Change --> Existing
Existing -->|Yes| Use
Existing -->|No| Add
Use --> Diagnose
Add --> Diagnose
Diagnose --> Optimize
```
Examples of future demand-driven areas may include:
- checkpoint coordinator fan-in;
- serialization or cross-task transfer;
- additional state/recovery paths;
- scheduler hot paths;
- long-running stability workloads.
They are not mandatory unfinished parts of STIP-32.
---
### 11. Future history and regression gate
Long-term result storage and a blocking regression gate remain optional
future evolution.
```mermaid
flowchart LR
Today["Current framework<br/>versioned result artifacts"]
History["Long-term storage"]
Trend["Performance history"]
Signal["Regression signal"]
Gate["Optional blocking gate"]
Today -.-> History
History --> Trend
Trend --> Signal
Signal -.-> Gate
```
The current result format already contains the data required to build this
later.
The missing requirement is mainly operational:
```text
controlled runner
stable environment
historical variance
rerun policy
ownership
false-positive control
```
Until those requirements exist, the benchmark framework should remain an
evidence and diagnosis tool rather than a blocking CI gate.
---
## Design summary
The final architecture follows a few simple rules:
1. **Measure production paths, not mocks.**
2. **Keep setup and validation outside the measurement unless intentionally
measured.**
3. **Validate correctness before trusting performance.**
4. **Use fixed logical work for persistent mutation benchmarks.**
5. **Keep the scheduled core suite small and stable.**
6. **Separate benchmark measurement from profiling diagnostics.**
7. **Do not change production APIs only for benchmark convenience.**
8. **Add new benchmarks when a real performance question requires them.**
The framework should now primarily be used to improve Zeta rather than
continuously expand benchmark coverage.
--
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]