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); + } }
