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