DanielLeens commented on PR #10874: URL: https://github.com/apache/seatunnel/pull/10874#issuecomment-5369108505
Thanks @SEZ9 — I rechecked `MultiTableSinkWriter.java` (and `MultiTableWriterRunnable.java`) at `95f1ee70765f` against your findings rather than taking them at face value, since I had approved this same head a few hours before your review landed. I need to walk that back: my last round only re-verified the two gaps I had previously raised (old-checkpoint restore undercounting, the newly-blocking enumerator→reader RPC) and did not re-audit the close-table logic inside `MultiTableSinkWriter` itself — which is exactly where your issues live. Confirmed against the current source: - Issue 1 (fast-path bypasses aggregation): real. `handleCloseTableEvent` (`MultiTableSinkWriter.java:580-585`) calls `markTablePendingClose` unconditionally whenever `sourceSubtaskId`/`expectedSourceEventCount` is null or `<=1`, with no check against `expectedCloseTableEventCounts`/`closeTableEventSources` for an already-established higher count from other readers. A single metadata-stripped event can close a table other readers are still writing to. - Issue 2/3 (drain race): also real. `hasQueuedRows` (`MultiTableSinkWriter.java:795-804`) only scans live `BlockingQueue` contents. In `MultiTableWriterRunnable.run()` (`MultiTableWriterRunnable.java:110-122`), `queue.poll(...)` removes the element from the queue *before* the worker enters `synchronized (this)` to process it. That's a genuine window where a row is already dequeued but not yet written, `hasQueuedRows` reports false, and `closeTable()`'s own `synchronized (runnable.get(i))` (`MultiTableSinkWriter.java:740`) doesn't close it — it only prevents a second concurrent close, not an already-in-flight dequeue that hasn't reached the monitor yet. - Issue 4 (leaked tracking maps) and Issue 6 (`Math::max` stale pin): both check out too. `handleCloseTableEvent` never short-circuits on `closedTableIds` before touching the maps (only `markTablePendingClose` does, and only after already having repopulated them), and the merge at line 590 is a plain `Math::max` with no path to correct a stale over-count from a reader with outdated state. - Issue 7 (`closedTableIds.add` before close succeeds): confirmed as well — `closeTable()` adds to `closedTableIds` and removes writers from `writerMap`/`sinkWriters` (lines 729, 755-756) before `sinkWriter.close()` is attempted, so a close failure leaves that writer un-retriable. I haven't re-verified 5/8/9/10 (docs, 2PC timing, javadoc) to the same depth, but nothing I checked contradicts them, and 9 in particular (closing inside `snapshotState()` before checkpoint commit) looks structurally accurate from the same read. This reopens the PR from my side — please treat my APPROVED review as superseded; I only have read access here so I can't dismiss it myself. Issues 1 and 2/3 are the ones I'd insist on before another look (both are real correctness gaps in the exact multi-reader/parallelism scenario this PR's own description calls out as the hard part), with 4/6/7 as high-value follow-ups in the same pass. Sorry for the back-and-forth, @davidzollo — I should have re-audited the sink-side close path from scratch instead of scoping my last round to only the two issues I'd personally flagged before. -- 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]
