DanielLeens opened a new issue, #12256:
URL: https://github.com/apache/seatunnel/issues/12256
### Search before asking
- [x] I searched the existing issues and pull requests.
- [x] This is a distinct residual path after #10189 and is not covered by
#12095.
### SeaTunnel Version
`dev` at `cf67b549a7a6c35fa0beb12d83c62892427ea919`
### Problem
`SinkAggregatedCommitterTask` retains one `checkpointBarrierCounter` entry
for every successful checkpoint whose sink writers return no commit information.
The cleanup added by #10189 removes barrier counters only while iterating
`checkpointCommitInfoMap`. Empty checkpoints never create an entry in that
sibling map, so their counters remain reachable for the lifetime of a
long-running streaming task.
With a five-second checkpoint interval, one affected task can retain 17,280
additional entries per day. Multiple long-running jobs multiply that growth and
can eventually increase GC pressure or exhaust heap.
### Reproduction
Reproduced on the `bigdata3` VM:
- Host: `bigdata3`
- Java: OpenJDK `1.8.0_502`
- SeaTunnel commit: `cf67b549a7a6c35fa0beb12d83c62892427ea919`
The regression test initializes a real `SinkAggregatedCommitterTask`,
records 10,000 barrier-counter entries without commit payloads, and invokes
`notifyCheckpointComplete(10000)`.
```text
[ERROR] testCheckpointBarrierCountersAreCleanedWithoutCommitInfo
completed empty checkpoints must not retain barrier counters
expected: <true> but was: <false>
Tests run: 1, Failures: 1
```
Command:
```bash
./mvnw -nsu -Dmaven.gitcommitid.skip=true \
-pl seatunnel-engine/seatunnel-engine-server \
-DskipITs -DskipIT=true \
-Dtest=SinkAggregatedCommitterTaskTest#testCheckpointBarrierCountersAreCleanedWithoutCommitInfo
\
test
```
### Root cause
1. `SinkFlowLifeCycle.processCheckpointBarrier()` calls
`writer.prepareCommit(checkpointId)`.
2. When it returns `Optional.empty()`, `SinkPrepareCommitOperation` carries
`commitInfos == null`.
3. `SinkPrepareCommitOperation.runInternal()` skips
`receivedWriterCommitInfo()`, so no `commitInfoCache` or
`checkpointCommitInfoMap` entry is produced.
4. `SinkAggregatedCommitterTask.triggerBarrier()` still inserts
`checkpointBarrierCounter[checkpointId]`.
5. `notifyCheckpointComplete()` removes counters only for keys found in
`checkpointCommitInfoMap`, so the empty checkpoint is never evicted.
### Expected behavior
Successful completion must remove all barrier-counter entries at or below
the completed checkpoint ID, regardless of whether commit information exists.
### Proposed fix
Decouple `checkpointBarrierCounter` eviction from `checkpointCommitInfoMap`
and add a regression test for many empty successful checkpoints. This does not
change checkpoint payloads, commit ordering, connector APIs, or serialized
state.
--
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]