DanielLeens commented on PR #12216: URL: https://github.com/apache/seatunnel/pull/12216#issuecomment-5846321542
Thanks for running this empirically rather than just reasoning about it in the abstract — the numbers line up exactly with what I'd expect from the bug this PR targets, and it's valuable to have that confirmed by someone who wasn't already anchored on my own review. On the `dev` vs. this PR vs. #12454 comparison: 8/14 failures on `dev`, matching the two-phase signature #12116 describes (first writer fails `MIN_ROW` on a partial total, then restarts fail `MAX_ROW` against never-reset static counters), against 0/13 on this PR's `AssertSinkWriter`, is consistent with the fix actually closing that window rather than just narrowing it. The 5/5 vs. 2/5 result for `AssertSinkWriterCloseTest` against the unpatched `AssertSinkWriter` is a good sanity check that the regression test is load-bearing rather than passing by construction. On your two questions about the "last open writer in this JVM" rule, I went back to the current head (`f424385`) to check both against the actual code rather than from memory: 1. You're right that `OPEN_WRITERS` counts writers by construction order, not by the sink's configured parallelism — the constructor does `OPEN_WRITERS.computeIfAbsent(openWritersKey, key -> new AtomicInteger()).incrementAndGet()` with no reference anywhere in the file to subtask count or parallelism. In the steady-state case this doesn't actually bite: every engine constructs all of a table's parallel writers together, before any of them has written enough rows to close, so the full population is always registered by the time the first `close()` runs. Where it could bite is a table whose writer set grows after a sibling has already closed — for example a dynamically-added table in a schema-evolving multi-table job. That's a real, if narrow, gap, and keying off the writer's actual subtask/parallelism count instead of construction order would remove the ambiguity outright. I don't think it changes the merge calculus for the flake this PR was written to fix, since that flake hap pens on the ordinary parallel-close path, not this one — but it's worth tracking as a real follow-up rather than dropping it. 2. This one isn't new on my end — it's the same Issue 2 I raised after the 2026-09-13 revision and reaffirmed as still open after CI ran on the current head on 2026-09-15: a writer that never reaches a graceful `close()` (task failure, forced kill, a retried attempt whose predecessor's slot was never released) leaves that table's `OPEN_WRITERS` count permanently above zero for the life of the JVM/classloader, and `MIN_ROW`/`MAX_ROW` silently stop being evaluated for every writer that closes after it. I still haven't been able to confirm from this connector's own code whether Zeta, Flink, or Spark guarantee `writer.close()` runs on task failure, which is exactly the question that needs answering before this can merge. A test for "writer dropped without a close, new writer created for the same table" is a good way to pin the actual behavior down either way — I'd support adding it. On #12454: if it's scoped to just the one `parallelism = 1` conf, then it only removes the specific reproduction under discussion here, not the underlying exposure for every other multi-table Assert conf running with `parallelism > 1`. Assuming that characterization holds, I agree this PR is the more complete fix of the two — conditional on point 2 above getting a real answer rather than staying open. For the record: `origin/dev` still carries the pre-fix `AssertSinkWriter` (no `OPEN_WRITERS`, each writer still evaluates on its own `close()`), so the underlying #12116 bug is not fixed upstream yet — this PR, or something like it, is still needed. -- 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]
