Rangsh opened a new pull request, #12081:
URL: https://github.com/apache/seatunnel/pull/12081
[Improve][Zeta] Reduce checkpoint state-store latency variance from
redundant WAL sync (#12058)
<!--
Thank you for contributing to SeaTunnel! Please make sure that your code
changes
are covered with tests. And in case of new features or big changes
remember to adjust the documentation.
Feel free to ping committers for the review!
## Contribution Checklist
- Make sure that the pull request corresponds to a [GITHUB
issue](https://github.com/apache/seatunnel/issues).
- Name the pull request in the form "[Feature] [component] Title of the
pull request", where *Feature* can be replaced by `Hotfix`, `Bug`, etc.
- Minor fixes should be named following this pattern: `[hotfix] [docs] Fix
typo in README.md doc`.
-->
### Purpose of this pull request
Fixes #12058.
Investigate and reduce the within-run latency variance observed on:
- `CheckpointStorageBenchmark.checkpointIdAtomicIncrement`
- `CheckpointStorageBenchmark.checkpointOverviewIncrementalUpdate`
Root cause:
- Both methods exercise the write-through IMap MapStore path
(`write-delay-seconds: 0`).
- Every measured operation waits for one file-backed WAL append.
- `HdfsWriter.flush()` previously performed redundant `hsync`/`hflush` calls
(up to three syncs on the HDFS path), which increased latency and CV without
improving durability.
- Sample-to-sample CV is dominated by durable WAL sync cost (plus overview
serialization size), not by the benchmark fixture and not by a Java 8 vs Java
11 comparison.
This PR applies a focused production-side fix while preserving checkpoint
correctness and state-store durability:
1. Keep exactly one durable sync path in `HdfsWriter.flush()`.
2. Fix `RequestFuture` completion/success semantics and make batch WAL waits
use the configured write timeout.
3. Ensure `WALWorkHandler` always publishes `done()` for APPEND events, even
on non-`IOException` failures.
4. Replace hot-path stream pipelines in overview stats /
`calculateStateSize` with simple loops to reduce GC noise on the measured path.
5. Document the variance source in `docs/en` and `docs/zh` benchmark guides,
and strengthen related unit coverage.
### Does this PR introduce _any_ user-facing change?
Yes, documentation only.
- Updated `docs/en/engines/zeta/benchmark.md` and
`docs/zh/engines/zeta/benchmark.md` to explain that the high CV on the two
checkpoint state-store microbenchmarks is dominated by write-through durable
WAL sync (and overview serialization size), and should not be treated as a
cross-JDK comparison signal.
- No user-facing config option, default value, public API, or checkpoint
recovery semantic change.
- Durability is preserved: each WAL append still completes only after one
durable `hsync`.
### How was this patch tested?
Unit tests:
- `RequestFutureTest`
- `CheckpointMonitorServiceCalculateStateSizeTest`
- `HazelcastCheckpointOverviewStateStoreTest`
- `HdfsWriterDurableFlushTest` (enabled on Linux/macOS, matching existing
imap-storage-file WAL tests; skipped on Windows when `HADOOP_HOME` is unset)
Local commands:
```bash
./mvnw -pl
seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file,seatunnel-engine/seatunnel-engine-server
spotless:apply
./mvnw -pl
seatunnel-engine/seatunnel-engine-storage/imap-storage-plugins/imap-storage-file
\
-Dtest=RequestFutureTest,HdfsWriterDurableFlushTest \
-DfailIfNoTests=false test
./mvnw -pl seatunnel-engine/seatunnel-engine-server \
-Dtest=CheckpointMonitorServiceCalculateStateSizeTest,HazelcastCheckpointOverviewStateStoreTest
\
-DfailIfNoTests=false test
```
Benchmark / profiling follow-up for reviewers (issue expected outcome):
- Use the same JDK, runner, JMH arguments, and state-store configuration for
before/after comparison.
- Suggested selectors:
- `CheckpointStorageBenchmark.checkpointIdAtomicIncrement$`
- `CheckpointStorageBenchmark.checkpointOverviewIncrementalUpdate$`
- Prefer the `Benchmarks Diagnostics` workflow on this PR branch vs `dev` to
collect JMH score + CV and optional CPU/wall/JFR evidence.
- Compare both operation latency (`us/op`) and within-run variance (CV /
Error) before and after this change.
### Check list
* [x] If any new Jar binary package adding in your PR, please add License
Notice according
[New License
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
* [x] If necessary, please update the documentation to describe the new
feature. https://github.com/apache/seatunnel/tree/dev/docs
* [x] If necessary, please update `incompatible-changes.md` to describe the
incompatibility caused by this PR.
* [x] If you are contributing the connector code, please check that the
following files are updated:
1. Update
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
and add new connector information in it
2. Update the pom file of
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
3. Add ci label in
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
4. Add e2e testcase in
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
5. Update connector
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
--
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]