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-12216-bf9749c191c4220b55ae5e05179fe1e71a3b1674
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit fea1ac0a4309d309a8b3c58fe746a2ae0243c006
Author: Daniel <[email protected]>
AuthorDate: Mon Oct 5 15:14:10 2026 +0000

    [Fix][Connector-V2][Assert] Evaluate row-count rules once per table after 
the last writer closes (#12216)
    
    Co-authored-by: DanielLeens <[email protected]>
    Co-authored-by: David Zollo <[email protected]>
---
 .../seatunnel/assertion/sink/AssertSinkWriter.java | 46 ++++++++++++
 .../assertion/sink/AssertSinkWriterCloseTest.java  | 84 ++++++++++++++++++++++
 2 files changed, 130 insertions(+)

diff --git 
a/seatunnel-connectors-v2/connector-assert/src/main/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriter.java
 
b/seatunnel-connectors-v2/connector-assert/src/main/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriter.java
index 0cb4a18adb..1c09ffab53 100644
--- 
a/seatunnel-connectors-v2/connector-assert/src/main/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriter.java
+++ 
b/seatunnel-connectors-v2/connector-assert/src/main/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriter.java
@@ -36,6 +36,7 @@ import java.util.Objects;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.CopyOnWriteArraySet;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.LongAccumulator;
 
 public class AssertSinkWriter extends AbstractSinkWriter<SeaTunnelRow, Void>
@@ -48,8 +49,25 @@ public class AssertSinkWriter extends 
AbstractSinkWriter<SeaTunnelRow, Void>
     private static final AssertExecutor ASSERT_EXECUTOR = new AssertExecutor();
     private static final Map<String, LongAccumulator> LONG_ACCUMULATOR = new 
ConcurrentHashMap<>();
     private static final Set<String> TABLE_NAMES = new CopyOnWriteArraySet<>();
+
+    /**
+     * Number of writers still open in this JVM, per table. The row counters 
above are shared by
+     * every parallel writer of a table (Flink subtasks, Zeta parallel tasks 
in one node), so the
+     * MIN_ROW / MAX_ROW rules describe the whole table and can only be judged 
once all of its
+     * writers have finished. The last writer of a table to close evaluates 
them; an earlier close
+     * would see a partial total and fail spuriously, which is the flaky 
behaviour reported for
+     * Flink in apache/seatunnel#12116.
+     */
+    private static final Map<String, AtomicInteger> OPEN_WRITERS = new 
ConcurrentHashMap<>();
+
     private final String catalogTableName;
 
+    /** Key of this writer in {@link #OPEN_WRITERS}; ConcurrentHashMap does 
not accept null. */
+    private final String openWritersKey;
+
+    /** Guards the open-writer count so a repeated close releases it exactly 
once. */
+    private boolean closed;
+
     public AssertSinkWriter(
             SeaTunnelRowType seaTunnelRowType,
             Map<String, List<AssertFieldRule>> assertFieldRules,
@@ -61,6 +79,8 @@ public class AssertSinkWriter extends 
AbstractSinkWriter<SeaTunnelRow, Void>
         this.assertRowRules = assertRowRules;
         this.assertTableRule = assertTableRule;
         this.catalogTableName = catalogTableName;
+        this.openWritersKey = catalogTableName == null ? "" : catalogTableName;
+        OPEN_WRITERS.computeIfAbsent(openWritersKey, key -> new 
AtomicInteger()).incrementAndGet();
     }
 
     @Override
@@ -102,6 +122,16 @@ public class AssertSinkWriter extends 
AbstractSinkWriter<SeaTunnelRow, Void>
 
     @Override
     public void close() {
+        if (closed) {
+            return;
+        }
+        closed = true;
+        if (!releaseAndCheckLastWriter()) {
+            // Another writer of this table is still running in this JVM and 
may still add rows,
+            // so the shared counters are not final yet. That writer evaluates 
the rules when it
+            // closes.
+            return;
+        }
         if (!assertRowRules.isEmpty()) {
             assertRowRules.entrySet().stream()
                     .filter(
@@ -170,4 +200,20 @@ public class AssertSinkWriter extends 
AbstractSinkWriter<SeaTunnelRow, Void>
                             + assertTableRule.getTableNames());
         }
     }
+
+    /**
+     * Releases this writer's slot in {@link #OPEN_WRITERS} and reports 
whether it was the last open
+     * writer of its table in this JVM. The entry is dropped at zero so a 
writer recreated later,
+     * for example after a restart, starts a fresh count instead of reusing a 
stale one.
+     *
+     * @return true when no other writer of the same table is still open
+     */
+    private boolean releaseAndCheckLastWriter() {
+        AtomicInteger openWriters = OPEN_WRITERS.get(openWritersKey);
+        if (openWriters == null || openWriters.decrementAndGet() > 0) {
+            return openWriters == null;
+        }
+        OPEN_WRITERS.remove(openWritersKey, openWriters);
+        return true;
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-assert/src/test/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriterCloseTest.java
 
b/seatunnel-connectors-v2/connector-assert/src/test/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriterCloseTest.java
index 75704dc410..e0c9df4e30 100644
--- 
a/seatunnel-connectors-v2/connector-assert/src/test/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriterCloseTest.java
+++ 
b/seatunnel-connectors-v2/connector-assert/src/test/java/org/apache/seatunnel/connectors/seatunnel/assertion/sink/AssertSinkWriterCloseTest.java
@@ -84,4 +84,88 @@ public class AssertSinkWriterCloseTest {
         writerA.write(new SeaTunnelRow(new Object[] {1}));
         Assertions.assertThrows(AssertConnectorException.class, 
writerA::close);
     }
+
+    /**
+     * Two parallel writers of one table share the row counter, so the MIN_ROW 
rule for the whole
+     * table must not fail when the first writer closes before the second one 
has written its rows.
+     * Only the last writer to close evaluates the rules, and it sees the full 
total.
+     */
+    @Test
+    public void testRowCountRulesAreEvaluatedByLastWriterOfTable() {
+        String tableName = "shared_min_row_table";
+        Map<String, List<AssertFieldRule.AssertRule>> assertRowRules =
+                Collections.singletonMap(tableName, 
Collections.singletonList(minRowRule(2)));
+
+        AssertSinkWriter firstWriter = newWriter(assertRowRules, tableName);
+        AssertSinkWriter secondWriter = newWriter(assertRowRules, tableName);
+
+        firstWriter.write(new SeaTunnelRow(new Object[] {1}));
+        Assertions.assertDoesNotThrow(
+                firstWriter::close,
+                "the first writer must not judge MIN_ROW while the second 
writer is still open");
+
+        secondWriter.write(new SeaTunnelRow(new Object[] {2}));
+        Assertions.assertDoesNotThrow(
+                secondWriter::close, "the last writer sees two rows in total, 
which meets MIN_ROW");
+    }
+
+    /**
+     * The rules describe the whole table, so the aggregated count of all 
parallel writers is what
+     * MAX_ROW is checked against: the last writer to close must fail when the 
total exceeds it.
+     */
+    @Test
+    public void testLastWriterFailsWhenAggregatedRowCountViolatesRule() {
+        String tableName = "shared_max_row_table";
+        Map<String, List<AssertFieldRule.AssertRule>> assertRowRules =
+                Collections.singletonMap(tableName, 
Collections.singletonList(maxRowRule(1)));
+
+        AssertSinkWriter firstWriter = newWriter(assertRowRules, tableName);
+        AssertSinkWriter secondWriter = newWriter(assertRowRules, tableName);
+
+        firstWriter.write(new SeaTunnelRow(new Object[] {1}));
+        Assertions.assertDoesNotThrow(firstWriter::close);
+
+        secondWriter.write(new SeaTunnelRow(new Object[] {2}));
+        Assertions.assertThrows(AssertConnectorException.class, 
secondWriter::close);
+    }
+
+    /**
+     * Closing a writer twice must release its open-writer slot only once; 
otherwise a repeated
+     * close could evaluate the rules while another writer of the table is 
still running.
+     */
+    @Test
+    public void testRepeatedCloseReleasesOpenWriterSlotOnce() {
+        String tableName = "shared_double_close_table";
+        Map<String, List<AssertFieldRule.AssertRule>> assertRowRules =
+                Collections.singletonMap(tableName, 
Collections.singletonList(minRowRule(3)));
+
+        AssertSinkWriter firstWriter = newWriter(assertRowRules, tableName);
+        AssertSinkWriter secondWriter = newWriter(assertRowRules, tableName);
+
+        firstWriter.write(new SeaTunnelRow(new Object[] {1}));
+        Assertions.assertDoesNotThrow(firstWriter::close);
+        Assertions.assertDoesNotThrow(
+                firstWriter::close, "a repeated close must not evaluate the 
rules early");
+
+        secondWriter.write(new SeaTunnelRow(new Object[] {2}));
+        // Two rows in total but MIN_ROW is three: the last writer still 
reports the violation.
+        Assertions.assertThrows(AssertConnectorException.class, 
secondWriter::close);
+    }
+
+    private static AssertFieldRule.AssertRule maxRowRule(int maxRows) {
+        AssertFieldRule.AssertRule rule = new AssertFieldRule.AssertRule();
+        rule.setRuleType(AssertFieldRule.AssertRuleType.MAX_ROW);
+        rule.setRuleValue((double) maxRows);
+        return rule;
+    }
+
+    private static AssertSinkWriter newWriter(
+            Map<String, List<AssertFieldRule.AssertRule>> assertRowRules, 
String tableName) {
+        return new AssertSinkWriter(
+                ROW_TYPE,
+                Collections.emptyMap(),
+                assertRowRules,
+                new AssertTableRule(Collections.emptyList()),
+                tableName);
+    }
 }

Reply via email to