DanielLeens opened a new pull request, #12216: URL: https://github.com/apache/seatunnel/pull/12216
## Purpose Fixes #12116: `TestFilterRowKindIT#testFilterRowKindMultiTable` fails intermittently on Flink-only CI legs (row counts like `75`/`97`/`125`/`150` instead of the expected `100`), which surfaces on many unrelated PRs since it is a pre-existing bug in `connector-assert`, not something any of those PRs' own diffs touch. Reproduced again this week, identically, on: - #12064 (Flink 11 leg) - #12150 (Flink 8 leg) - #12182 (Flink 11 leg) ...and previously on #11077/#11458/#11503 (see #12116 for full evidence). ## Root cause `AssertSinkWriter` keeps its row-count and table-name rule state in two `static` maps (`LONG_ACCUMULATOR`, `TABLE_NAMES`), keyed only by table name, and neither is ever cleared: - On Zeta, each job runs under its own classloader, so the statics are effectively scoped per job. - On Flink, every per-table writer and every sink subtask across the whole TaskManager JVM shares those same statics, and each writer's `close()` evaluates its own MIN_ROW/MAX_ROW or table-set rule against them independently. Since a table name can be (and in this suite's case, is) reused across unrelated jobs/tests within one JVM lifetime, and nothing ever removes an entry: - a writer that closes before a same-named table in a *different*, still-running job finishes contributing can see a partial total (`MIN_ROW` failures - the `75`/`97` cases), and - a writer whose table name collides with an already-accumulated (and never cleared) entry from an *earlier* job can see an inflated total (`MAX_ROW` failures - the `125`/`150` cases, i.e. 100 + 25/50 leftover). Both directions appearing across the same suite's runs matches shared mutable static state plus per-writer independent evaluation, not test data or a race in the transform/source under test. ## Fix Give each `AssertSink` construction (one per job/table) a fresh `sinkInstanceId`, and thread it into every `AssertSinkWriter` that sink creates. Both static maps are now keyed by `(sinkInstanceId, tableName)` / `sinkInstanceId` instead of table name alone. This keeps the *intentional* behavior the statics exist for - aggregating row counts and table names across every parallel subtask writer of the **same** sink instance, which Flink's shared-JVM subtask execution model still needs - while eliminating cross-job/cross-test contamination. Not a test-only workaround: the test and its assertion are already correct (see #12116); this fixes the shared production connector code whose bug the test's failures were exposing. Minimal, contained change: no new constructor parameters on any public-facing `Sink`/`SinkFactory` API, no change to `AssertSink`'s own public constructor - only `AssertSinkWriter`'s package-internal construction gained one new parameter, updated at its one production call site (`AssertSink#createWriter`) and its two existing direct-construction unit tests. ## Tests Added three tests to `AssertSinkWriterCloseTest`: - `testRowCountIsIsolatedAcrossUnrelatedSinkInstances` - two sink instances (different `sinkInstanceId`, same table name) each write 1 row against `MAX_ROW=1`; both must close successfully, proving one job's count can no longer leak into the other's. - `testRowCountStillAggregatesAcrossSubtasksOfSameSinkInstance` - two writers sharing the *same* `sinkInstanceId` and table (simulating two parallel subtasks), writing 3 and 2 rows respectively against `MIN_ROW=5`; the second writer's `close()` must see the combined total of 5, proving the intentional cross-subtask aggregation still works. - `testTableNamesRuleIsIsolatedAcrossUnrelatedSinkInstances` - same isolation proof for the `assert_table` table-set rule. Also updated the two existing `AssertSinkWriterCloseTest` cases to pass a `sinkInstanceId` (previously not needed to reproduce their scenario) - their existing assertions are unchanged. Per this repository's contribution guidance for Apache SeaTunnel, local execution was limited to `./mvnw spotless:apply`/`spotless:check` (passing); full compilation, the new/existing unit tests, and any E2E impact are verified by this PR's own GitHub Actions CI run. -- 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]
