DanielLeens commented on PR #12081:
URL: https://github.com/apache/seatunnel/pull/12081#issuecomment-5611997422
@SEZ9 direct answers below, re-verified against `64b0ec85a` file-by-file
rather than trusting my prior summary.
**On the "cut off" text:** pulled my `09-09T11:48` comment fresh via the API
— it is not actually truncated. The sentence after `This is exactly the `
continues: `"stop cleanly at an incomplete trailing frame, keep everything read
so far" behavior the comment claims — confirmed by reading the reader's source,
not by trusting the comment.` I'd guess a client-side rendering quirk with the
literal `"` right after "the " (you reported an identical "looks cut off
mid-sentence" symptom on another PR today too, so this may be reproducible on
your client rather than anything actually missing here).
**On F4 (mid-file torn record) — you're right, I'm reopening this rather
than standing by my "resolved" call.** Traced `WALWorkHandler.walEvent()`
(`WALWorkHandler.java:71-84`): after a caught write failure it does **not**
reset/reopen `writer`, and the very next `APPEND` event still calls
`writer.write()` on the same still-open `HdfsWriter`/`FSDataOutputStream` — the
worker keeps consuming events on the same writer for the rest of the process's
life, it never stops after a failure. So "the torn record stays the file's last
byte" is not actually guaranteed by this code; it only holds if nothing after
it is ever written successfully, which the current design doesn't enforce.
Whether that's exploitable depends on which half of
`HdfsWriter.write(byte[])` fails (`HdfsWriter.java:82-86`, calls
`out.write(data)` then `this.flush()`):
- If only `flush()` (the `hsync`) throws, `out.write(data)` already fully
placed the record's bytes into the stream — nothing is torn, just not yet
confirmed durable, and the *next* successful record's `hsync` would flush it
along too.
- If `out.write(data)` itself throws partway through the byte array, the
file can end up with a genuinely truncated record followed by a later,
fully-successful one. `DefaultReader.readData()`'s loop breaks (not continues)
the instant a length prefix claims more bytes than remain
(`DefaultReader.java:100-104`), so that later, otherwise-intact record — and
everything appended after it, forever — would be silently lost on the next
read, not just the one failed record.
The in-code comment at `WALWorkHandler.java:76-78` ("HdfsWriter/CloudWriter
serialize before mutating the stream, so unchecked failures ... do not leave a
torn mid-file record") asserts the second bullet can't happen, but doesn't show
why `FSDataOutputStream.write(byte[])` can't itself partially advance before
throwing — and I don't have a source-level answer for that either; it depends
on internal DFSOutputStream/DataStreamer packet-flush behavior on a large
single `write()` call, which I haven't traced. Given that gap, I agree with
keeping F4 open until either the failure mode is pinned down, or
`DefaultReaderTornTrailingRecordTest` is extended to cover "torn record with a
valid record written after it" and the observed behavior (data loss vs. safe
skip-and-continue) is made explicit.
**Quick status on the rest, verified directly against `64b0ec85a`:**
- **F1/F3** (Future contract — `isDone()` no longer conflated with success;
timed `get()` throws `TimeoutException` instead of returning `false`) —
resolved. `RequestFuture.java:31-32` (class Javadoc), `84-89` (`get(timeout,
unit)` throws).
- **F2** (`WALWorkHandler` catches `Exception`, `executeResponse()`
separately guarded) — resolved. `WALWorkHandler.java:73` (catch-all), `98-108`
(own try/catch around `executeResponse`).
- **F5** (Mockito dependency) — resolved. `imap-storage-file/pom.xml` has an
explicit `mockito-junit-jupiter` test-scope dependency with a comment
explaining it's kept local even though also inherited from the parent — not
relying on an implicit transitive dep.
- **F6** (shared batch deadline) — resolved. `IMapFileStorage.java:346-347`
computes `deadlineNanos` once before the loop; `352-353` clamps remaining time
to `Math.max(0L, ...)` and skips the wait entirely at zero — so the whole batch
is bounded by one `writDataTimeoutMilliseconds`, not N×.
- **F7** (`get()`/`get(timeout, unit)` Javadoc) — resolved. Both now
individually documented: `RequestFuture.java:57-63` (untimed, states no
production caller), `70-83` (timed, `@throws TimeoutException`).
- **F8** (timeout logging) — resolved in both call sites.
`queryExecuteStatus` (`IMapFileStorage.java:324-331`) and
`batchQueryExecuteFailsStatus` (`IMapFileStorage.java:366-375`) both log one
`WARN` with requestId/elapsed/limit and push the stack trace to `DEBUG`,
instead of an `ERROR`-level stack per timeout.
So F1/F2/F3/F5/F6/F7/F8 all check out cleanly against the current head —
only F4 is going back on the open list, and only for the specific
"torn-then-more-writes" gap above.
--
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]