davidzollo commented on PR #10874:
URL: https://github.com/apache/seatunnel/pull/10874#issuecomment-5579089535

   Pushed `ce9c17406b7` addressing F5(comment)/F5-F7/F1-F8/F4-F6; F3 is 
answered below rather than folded into the commit.
   
   - **F4/F6 (leak)** — `handleCloseTableEvent` now returns immediately if 
`closedTableIds.contains(event.tableId())`, before touching 
`closeTableEventSources`/`expectedCloseTableEventCounts`. A straggler event on 
an already-closed table can no longer recreate those entries, closing the leak: 
added `testStragglerCloseEventAfterCloseDoesNotLeakAggregationState`, which 
reads both maps via reflection and asserts they stay empty after a straggler 
event.
   - **F5 (comment)** — added the documentation you asked for right at the 
fall-through in the un-attributed-event branch (`markTablePendingClose` when 
`requiredCountOnFile` is null/`<=1`), spelling out that this is a known, 
narrower-but-not-eliminated ordering-ambiguity gap rather than an oversight.
   - **F5/F7 (late row bypasses policy)** — the already-closed-table check in 
`write()` now checks `failurePolicy.continueOtherTables()` before throwing: 
under continue policy it logs and drops the row (same outcome as the 
quarantined-table and missing-primary-key paths just below it); the default 
fail-fast policy is unchanged and still throws (covered by the existing 
`testCloseTableEventAllowsInFlightRowsUntilFinalSnapshot`). Added 
`testCloseTableWriteContinuesOtherTablesInsteadOfThrowing` for the 
continue-policy case.
   - **F1/F8 (unbounded wait)** — `waitUntilTableQueueDrained` now has a 300s 
deadline; hitting it throws `IOException` (the table stays in 
`pendingCloseTableIds` since this throws before `closeTable` removes it, so the 
next checkpoint retries the close instead of it being silently skipped). The 
`InterruptedException` branch now also throws `IOException` — matching the 
method's declared contract — instead of the `RuntimeException` wrap. I didn't 
add a dedicated test that actually waits out the deadline (would make the suite 
slower for a one-line comparison); the existing close tests keep exercising the 
normal-drain path through this method.
   - **F2** — left as the non-blocking follow-up you already tracked, unchanged 
in this pass.
   - **F3 (in-flight row between `queue.poll()` and write)** — traced this 
through `MultiTableWriterRunnable.run()`: `queue.poll()` happens outside any 
lock, but the actual `queueElement.process(this)` (and therefore `writeRow`'s 
`tableIdWriterMap.get(...)` lookup) happens inside `synchronized (this)` — and 
`closeTable()`'s writer removal (`sinkWriter.close()` + `writerMap.remove(...)` 
+ `runnable.get(i).removeTableWriter(tableId)`) happens inside `synchronized 
(runnable.get(i))`, the *same* monitor. So the two can't interleave 
arbitrarily; what can happen is: a row is polled off the queue (so 
`hasQueuedRows(tableId)` in `waitUntilTableQueueDrained` no longer sees it) 
just before `closeTable` acquires the lock and removes the writer, and then the 
worker thread acquires the lock second and calls `writeRow` against a writer 
map that's already had the entry removed. 
`tableIdWriterMap.get(row.getTableId())` returns null, 
`allowSingleWriterFallback` is false for a genuine multi-table 
 job, and (outside `continueOnTableFailure`) it throws 
`RuntimeException("...can't find writer for tableId...")`, which surfaces as a 
task failure rather than a silently dropped/lost row. So this is real, but it's 
a narrow, checkpoint-recoverable spurious-failure window, not silent data loss.
   
     A precise fix needs a **per-table** in-flight counter (the existing 
`pendingRowRequests` in `MultiTableWriterRunnable` is per-runnable/whole-queue, 
which would over-serialize unrelated tables sharing a queue if reused directly 
for this), incremented in `MultiTableSinkWriter.offerRowElement` and 
decremented once `MultiTableWriterRunnable.run()`'s `finally` block finishes 
processing that row — i.e. new shared state and a new decrement call spanning 
both classes, not a one-line change. I'd rather scope that as its own follow-up 
PR with its own review than fold a new cross-class counter into this one. Happy 
to take it on next if you'd like it landed before merge rather than after.
   
   `./mvnw spotless:apply` run on `seatunnel-api` before pushing, no further 
formatting changes produced beyond what's in this diff. Per the SeaTunnel 
local-verification rule I'm not running the module's tests locally — GitHub CI 
on this head is the verification source of truth.
   


-- 
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