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]

Reply via email to