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]

Reply via email to