This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 46e58fc16f [Fix][Zeta] Clean empty checkpoint barrier counters (#12260)
46e58fc16f is described below

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

Reply via email to