This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12260-35b2716cde7d4c91a24fc618a8d9cae90e213db3 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 46e58fc16f8523c82edd40b93e8f8d62663737a2 Author: Jast <[email protected]> AuthorDate: Mon Sep 14 03:03:38 2026 +0000 [Fix][Zeta] Clean empty checkpoint barrier counters (#12260) Co-authored-by: zhangshenghang <[email protected]> --- .../server/task/SinkAggregatedCommitterTask.java | 2 +- .../server/task/SinkAggregatedCommitterTaskTest.java | 20 ++++++++++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTask.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTask.java index 2a9a24d401..e2bf29371a 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTask.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTask.java @@ -303,6 +303,7 @@ public class SinkAggregatedCommitterTask<CommandInfoT, AggregatedCommitInfoT> @Override public void notifyCheckpointComplete(long checkpointId) throws Exception { List<AggregatedCommitInfoT> aggregatedCommitInfo = new ArrayList<>(); + checkpointBarrierCounter.keySet().removeIf(key -> key <= checkpointId); checkpointCommitInfoMap.forEach( (key, value) -> { if (key > checkpointId) { @@ -311,7 +312,6 @@ public class SinkAggregatedCommitterTask<CommandInfoT, AggregatedCommitInfoT> aggregatedCommitInfo.addAll(value); checkpointCommitInfoMap.remove(key); commitInfoCache.remove(key); - checkpointBarrierCounter.remove(key); }); List<AggregatedCommitInfoT> commit = aggregatedCommitter.commit(aggregatedCommitInfo); tryClose(checkpointId); diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTaskTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTaskTest.java index 434b291dce..2e7bee624c 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTaskTest.java +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/SinkAggregatedCommitterTaskTest.java @@ -134,6 +134,26 @@ public class SinkAggregatedCommitterTaskTest { "checkpointCommitInfoMap should still contain checkpoint 3"); } + @Test + void testCheckpointBarrierCountersAreCleanedWithoutCommitInfo() throws Exception { + Map<Long, Integer> checkpointBarrierCounter = getCheckpointBarrierCounter(); + checkpointBarrierCounter.put(1L, 1); + checkpointBarrierCounter.put(2L, 1); + checkpointBarrierCounter.put(3L, 1); + + task.notifyCheckpointComplete(2L); + + Assertions.assertFalse( + checkpointBarrierCounter.containsKey(1L), + "completed empty checkpoints must not retain barrier counters"); + Assertions.assertFalse( + checkpointBarrierCounter.containsKey(2L), + "completed empty checkpoints must not retain barrier counters"); + Assertions.assertTrue( + checkpointBarrierCounter.containsKey(3L), + "future checkpoints must retain their barrier counters"); + } + @Test void testCheckpointCacheCleanupAfterNotifyCheckpointAborted() throws Exception { // Simulate receiving commit info for a checkpoint
