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

   Thanks @SEZ9 — I pulled the raw review body back via the API to check, and 
it is not actually truncated: the full text runs through all five sections 
(Core Logic Analysis, Compatibility, Performance, Issue Summary, Merge 
Recommendation), ending with the "Overall assessment" paragraph. This looks 
like a GitHub rendering/collapse artifact on a long review body rather than a 
real cutoff on my end — sorry for the confusion either way. To make sure 
nothing is lost, here is the missing part inline:
   
   **F1/F8 (unbounded wait in `snapshotState`) — Fixed.** 
`waitUntilTableQueueDrained` (`MultiTableSinkWriter.java:848-874`) now has a 
300s deadline (`TABLE_QUEUE_DRAIN_TIMEOUT_MILLIS`), throws `IOException` on 
timeout with the table left in `pendingCloseTableIds` (verified the timeout 
throw happens before `closeTable`'s removal of that entry, at `:770-774`, so 
the next checkpoint retries), and the `InterruptedException` branch now 
restores the interrupt flag and throws `IOException` instead of an undeclared 
`RuntimeException`.
   
   **F4/F6 (straggler-event map leak) — Fixed.** `handleCloseTableEvent` 
(`:578-588`) now returns immediately when 
`closedTableIds.contains(event.tableId())`, before touching 
`closeTableEventSources`/`expectedCloseTableEventCounts`, so a late/duplicate 
event for an already-closed table can't resurrect those tracking-map entries. 
`testStragglerCloseEventAfterCloseDoesNotLeakAggregationState` asserts both 
maps stay empty via reflection.
   
   **F5/F7 (late row bypasses `MultiTableFailurePolicy`) — Fixed**, as you 
already confirmed: `:677-691` checks `failurePolicy.continueOtherTables()` 
before throwing, mirroring `:696-699`/`:706-714`. I'd treat this the same way 
you propose — resolved pending the continue-policy test, which exists 
(`testCloseTableWriteContinuesOtherTablesInsteadOfThrowing`).
   
   **F2 (no final commit before `closeTable()` closes a 2PC sub-writer) — Not 
fixed, tracked.** Still open at `:769-838`: `sinkWriter.close()` is called 
directly with no interposed `prepareCommit`/flush, so buffered-but-uncommitted 
data for that table is dropped at close time. I'm not blocking on this being 
code-fixed in this PR, but I do want it written into the PR description as a 
known, tracked limitation before merge, since right now it only lives in 
review-comment history.
   
   **F3 (queue-poll vs. `closeTable()` writer-removal race) — Not fixed, 
explained, deferred.** I traced the synchronization myself: `queue.poll()` 
(`MultiTableWriterRunnable.java:117`) is outside any lock, but 
`queueElement.process(this)`/`writeRow` is inside `synchronized (this)`, and 
`closeTable()`'s writer removal is inside `synchronized (runnable.get(i))` — 
the same monitor. A row dequeued just before `closeTable` removes its writer 
can still lose that race and hit a null `tableIdWriterMap.get(...)`, throwing 
under fail-fast policy. It's real but narrow, and it surfaces as a 
checkpoint-recoverable task failure, not silent data loss. Agreed a per-table 
in-flight counter spanning both classes is its own follow-up, not a one-line 
fix — but per my review, please file it as a tracked issue linked from this PR 
before merge, not left only in review history.
   
   Separately, while re-checking the current diff to answer this: 
`expectedCloseTableEventCounts.merge(event.tableId(), expectedSourceEventCount, 
Math::max)` in `handleCloseTableEvent` is still present at this head unchanged. 
That's the same merge strategy behind the stale-over-count-pin concern raised 
in earlier rounds (a stale, larger `expectedSourceEventCount` from an older 
reader generation can pin `requiredCount` above what will ever arrive, so the 
table never closes and just silently falls back to close-at-task-end with no 
warning log). F4/F6 above fixed the straggler-leak half of that area but is a 
distinct code path from the `Math::max` pin itself, and I don't see it called 
out in my own Issue Summary in the last full review — flagging so it isn't 
accidentally dropped from tracking. I don't have a fresh repro confirming it's 
still live given how much else changed around it, so I'm not reopening it as a 
numbered issue this round without doing that work properly — I'll give
  it a full re-trace in the next full review pass rather than assert its status 
here on a reply-only round.
   
   Bar for merge from my side stays: F2 written into the PR description 
(blocker), F3 filed as a tracked follow-up issue (recommended, non-blocking), 
and the CI picture cleaned up (sync + rerun, since the current failures are all 
in modules this PR doesn't touch). Once those land I don't expect further 
structural concerns from me.


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