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]

Reply via email to