davidzollo opened a new pull request, #12095:
URL: https://github.com/apache/seatunnel/pull/12095

   ## Purpose
   
   Adds multi-node E2E regression coverage for the checkpoint-cache memory leak
   fixed in #10189 ("[Fix][Zeta] Fix memory leak in 
SinkAggregatedCommitterTask"),
   reported in #10188.
   
   The original report: a production deployment running 25 concurrent MySQL-CDC-
   to-Iceberg streaming jobs, each with a 5-second checkpoint interval, saw heap
   usage grow from near-zero to approximately 18.6 GiB over about 3 days, with
   GC cycle count and duration climbing steadily throughout, before eventually
   crashing with `OutOfMemoryError`. Heap dump analysis traced the growth to two
   maps in `SinkAggregatedCommitterTask` — `commitInfoCache` and
   `checkpointBarrierCounter`, both keyed by checkpoint id — accumulating one
   stale entry per checkpoint cycle, forever, because `notifyCheckpointComplete`
   and `notifyCheckpointAborted` removed the completed/aborted checkpoint's
   entry from the sibling `checkpointCommitInfoMap` but never from these two.
   
   ## What was still missing
   
   The fix (two `.remove(key)` calls) already ships with a thorough single-JVM
   unit test (`SinkAggregatedCommitterTaskTest`) that drives a
   directly-instantiated task through mocked checkpoint callbacks and asserts
   the maps empty out. What that test cannot prove is the actual production
   scenario: a real streaming job, on a real cluster, running through many real
   checkpoint cycles back to back, staying bounded rather than accumulating.
   
   ## What this test does
   
   
`SinkAggregatedCommitterTaskCheckpointCacheLeakIT#testCommitInfoCacheStaysBoundedAcrossManyCheckpoints`:
   
   1. Starts a 1-master + 1-worker cluster and submits a real streaming job
      (`FakeSource` → `LocalFile` sink with transactions enabled, which routes
      through `FileSinkAggregatedCommitter`) at a 1-second checkpoint interval.
   2. Locates the live `SinkAggregatedCommitterTask` instance on the worker by
      scanning `TaskExecutionService`'s own execution-context bookkeeping (no
      public accessor exists from a client-computable `TaskGroupLocation`, since
      the committer task's own task id is an engine-internal implementation
      detail).
   3. Samples `commitInfoCache` and `checkpointBarrierCounter` sizes by
      reflection — the identical field-access technique
      `SinkAggregatedCommitterTaskTest` already uses, just against a live task
      instead of a directly-instantiated one — 20 times at a 500ms interval
      (spanning roughly 10 real checkpoint cycles at the job's 1-second
      interval).
   4. Asserts every single sample stays at or below a small fixed bound (3
      entries). A pre-fix regression would grow both maps roughly linearly with
      the checkpoint count, reaching close to 10 stale entries by the last
      sample — clearly distinguishable from "stays bounded," and distinguishable
      from "leaking slowly" in a way a single fixed-duration wait could not be.
   
   ## Test plan
   
   - New test and new test-resource config only; no production code changed.
   - `./mvnw spotless:apply` run on the affected module — succeeded.
   - `./mvnw install -DskipTests` (module + dependency chain, 
`seatunnel-engine-ui`
     excluded since it's unmodified) run to confirm the new test genuinely
     compiles — confirmed via `BUILD SUCCESS` and the actual `.class` file
     present under `target/test-classes` (not run with `-Dmaven.test.skip=true`,
     which skips test compilation entirely rather than just execution).
   - Full test execution is left to CI per this repository's E2E conventions
     (Hazelcast-cluster-backed, not run in this sandbox).


-- 
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