davidzollo opened a new issue, #12116:
URL: https://github.com/apache/seatunnel/issues/12116

   ### Search before asking
   - [x] I had searched in the issues and found no similar issues.
   
   ### What happened
   `TestFilterRowKindIT#testFilterRowKindMultiTable` 
(`filter_row_kind_exclude_insert_multi_table.conf`: FakeSource with 3 tables × 
100 rows, parallelism 1, BATCH; `FilterRowKind` drops INSERT for 
`test.abc`/`test.xyz`; Assert expects `abc=0, xyz=0, www=100`) fails 
intermittently **only on Flink legs** — Zeta and Spark pass in the same runs — 
with the Assert sink reporting wrong row counts for `test.www`:
   
   - apache/seatunnel PR #11458, fork run `zhangshenghang/seatunnel` 
33971790374, job `transform-v2-it-part-1 (11)`, Flink **1.18.0** leg: `row num 
:97 fail MIN_ROW 100`, `row num :125 / :150 fail MAX_ROW 100`
   - PR #11077, fork run `hesam-oxe/seatunnel` 33971836407, same job, Flink 
**1.15.3** leg: `row num :75 fail MIN_ROW 100`, `row num :125 / :150 fail 
MAX_ROW 100`
   - PR #11503, fork run `nielifeng/seatunnel` 33941523842, same job, same test 
(`expected: <0> but was: <1>` on the exit code)
   
   None of those PRs touch connector-assert, the transform, or the e2e module. 
All three heads already contain #11995 ("Remove fixed sleep in AssertSinkWriter 
close and assert own table only", merged 2026-09-01, `fa001a4e`), so that fix 
does not cover this.
   
   ### Root cause (from source, `dev` `AssertSinkWriter`)
   - Row counters are **static, JVM-wide**: `private static final Map<String, 
LongAccumulator> LONG_ACCUMULATOR` and `private static final Set<String> 
TABLE_NAMES` (lines 49–50).
   - `write()` accumulates into `LONG_ACCUMULATOR[element.getTableId()]` 
(multi-rule config path), and each writer's `close()` reads that static map for 
its own `catalogTableName` and throws `RULE_VALIDATION_FAILED` on 
MIN_ROW/MAX_ROW (lines 104–160).
   - On **Zeta**, each job runs with its own classloader, so the statics are 
effectively per job and per writer set. On **Flink**, every per-table writer 
and every sink subtask lives in the same TaskManager JVM and shares those 
statics, and each writer's `close()` evaluates independently:
     - a writer that closes before the other subtasks/writers finish 
accumulating sees a **partial** total → `MIN_ROW` failures (`75`, `97`);
     - accumulation from other writers/attempts under the same key inflates the 
total → `MAX_ROW` failures (`125`, `150` — i.e. 100 + 25/50 chunks).
   
   Both directions appear in the same job, which matches shared mutable static 
state plus per-subtask close(), not test data.
   
   ### What you expected
   Row-count rules evaluated against the rows the writer (or the job) actually 
received, isolated per job/table/writer, on every engine.
   
   ### Suggested fix (connector-assert, `src/main`)
   Make the counters instance/per-job state (e.g. per-writer counters 
aggregated at commit/close via the sink's commit path, or keyed by job id + 
table), and evaluate `MIN_ROW`/`MAX_ROW` once all writers for the table have 
finished rather than per-subtask `close()`. Not proposing a test-only 
workaround: the test is correct and the assertion is the intended contract.
   
   ### SeaTunnel Version
   dev (`04fabece`), reproduced on Flink 1.15.3 and 1.18.0 e2e containers.
   
   ### Engine
   Flink (Zeta and Spark unaffected).
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to