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]