[Improvement-18459][Common] Stream task log download in bounded chunks to prevent OOM - #18463
[Improvement-18459][Common] Stream task log download in bounded chunks to prevent OOM#18463xmg333 wants to merge 2 commits into
Conversation
SbloodyS
left a comment
There was a problem hiding this comment.
For a log larger than 47 MB:
LogServiceImpl#getTaskInstanceWholeLogFileBytesreturnsERROR.LogClientDelegate#getWholeLogBytestreats every local error as a reason to callremoteLogClient.getWholeLog(...).RemoteLogClientcallsgetFileContentBytesFromRemote, which now uses the same capped reader and silently returns only the first 47 MB.
With remote logging enabled, the API can therefore return a successfully downloaded but truncated log. With remote logging disabled or unavailable, it may return an empty log or a generic download error instead of the clear size-limit message.
Please distinguish “local log unavailable” from “log exceeds the supported size,” propagate the latter to the API, and make the reader fail explicitly rather than silently truncating. The regression test should cover the complete LogClientDelegate/API path, not only LogServiceImpl.
Additionally, the linked issue expects large logs to remain downloadable through chunked streaming. This PR rejects them entirely.
2aa5e28 to
f7168cf
Compare
|
Thanks for the review @SbloodyS . I've reworked the approach based on your feedback. New ** chunked streaming log ** is now completed. Could you confirm if this scope is what you had in mind? What's newNew RPC: Fallback logic (the important part)The API server runs a chunked loop with an The core invariant: fallback only happens when Three concrete scenarios:
Legacy path unchanged:
Tests cover all three scenarios in |
93169b9 to
b2731e8
Compare
| final byte[] bytes = remoteLogClient.getWholeLog(taskInstance); | ||
| if (bytes != null && bytes.length > 0) { | ||
| outputStream.write(bytes); | ||
| } | ||
| outputStream.flush(); |
There was a problem hiding this comment.
If remote returns null/empty (archive missing), this still flush()es and the download ends as HTTP 200 with only the log header. Please throw when bytes are absent so a missing remote log is not reported as a successful download.
There was a problem hiding this comment.
Thanks for catching this. Fixed in LogClientDelegate.writeRemoteLegacy: when remoteLogClient.getWholeLog(...) returns null/empty (remote archive missing), it now throws IOException instead of flushing an empty body. The exception propagates through streamWholeLog → StreamingResponseBody and aborts the response, so a missing log is no longer reported as a successful HTTP 200 download.
While reviewing the fix I also found and fixed a related gap: LoggerServiceImpl.checkDownloadLogAuth validated host but not logPath. It's now fixed and would return a clear error before streaming starts.
b2731e8 to
2191ccc
Compare
SbloodyS
left a comment
There was a problem hiding this comment.
Keep every fallback path memory-bounded
LogClientDelegate.java:132-202
The chunked path is bounded, but both fallbacks still load the complete log into memory:
- An old worker or first-chunk failure calls
localLogClient.getWholeLog(), whose worker-side implementation usesgetFileContentBytesFromLocal()and aByteArrayOutputStream. - A missing worker calls
remoteLogClient.getWholeLog(), which usesgetFileContentBytesFromRemote()and then reads the downloaded file into another completebyte[].
Therefore, downloading a large log from an old worker or remote storage can still reproduce the original OOM. Please stream the remote file from disk and either reject the unsupported old-worker path explicitly or otherwise make it enforceably bounded. Add regression coverage for large fallback logs.
Do not execute the whole-file fallback twice after an error
LogClientDelegate.java:141-164
When the first chunk returns a non-success response, writeLocalLegacy() is called inside the outer try. If that fallback throws—for example, the legacy RPC fails and the remote archive is missing—the outer catch still sees offset == 0 and calls writeLocalLegacy() again.
This repeats the legacy and remote requests and can also retry after a fallback has partially written to the response. Please limit the catch to the chunk RPC itself or otherwise let fallback failures propagate without re-entering the fallback.
BTW, the PR title and description still describe a 47 MB cap, while the implementation now uses chunked streaming. Please update them to match the current approach.
2191ccc to
2ecb4eb
Compare
|
Thanks for the three points — all have been addressed. 1. "Keep every fallback path memory-bounded"All three paths are now memory-bounded:
Additionally, Without this fix, the original exception cause could be lost and replaced by a Test coverage:
2. Do not execute the whole-file fallback twice after an error
This ensures that:
This is covered by:
The test explicitly verifies:
3. PR title and descriptionUpdated to: Stream task log download in bounded chunks to prevent OOM |
SbloodyS
left a comment
There was a problem hiding this comment.
NettyClientHandler stores only one opaque request ID in the channel attribute OPAQUE_KEY, while NettyRemotingClient reuses the same channel for concurrent RPC requests.
This creates a race:
- Request A stores opaque A.
- Request B stores opaque B on the same channel.
- A frame/decoder error occurs.
exceptionCaught()only completes opaque B.- Request A remains pending until its timeout.
There is also a second race: when any request completes successfully, doSendSync() unconditionally clears OPAQUE_KEY, which can erase the opaque ID of another request that is still in flight.
Please avoid using a single channel attribute for pending request tracking. On channel failure, all pending ResponseFutures associated with that channel should be completed with the original exception, or the channel should maintain a proper set/map of in-flight opaque IDs. Please also add a regression test with at least two concurrent requests sharing one channel and a decoder exception.
2ecb4eb to
5ac4d26
Compare
|
Thanks for the review. I spent quite a bit of time debugging the previous The new approach is simpler: each The lifecycle is now:
This also fixes a few cases that were easy to miss with the previous approach:
The requested regression test is included: It sends two concurrent requests over the same real Netty channel, triggers a malformed frame, and verifies that both callers receive the decoder error promptly. I also added Other fixesWhile testing this, I found a few highly related issues and fixed them as well:
VerificationThere are now 80 tests across the 4 modules, including regression tests for the cases above. |
SbloodyS
left a comment
There was a problem hiding this comment.
The legacy fallback is still unbounded on the worker side
The maxFrameSize check only protects the receiving API process. It does not make the legacy RPC memory-bounded on an old worker.
When the chunk RPC is unavailable during a rolling upgrade, the current flow is:
streamWholeLog()falls back towriteLocalLegacy().localLogClient.getWholeLog()invokes the old worker RPC.- The old worker reads the entire file into a
ByteArrayOutputStream. - It creates another full
byte[]withtoByteArray(). - The RPC layer JSON/base64-serializes the complete response.
- Only after all of those allocations does the API-side
TransporterDecoderreceive a frame and enforcemaxFrameSize.
Therefore, a sufficiently large log can still OOM the old worker before the decoder has anything to reject. JSON/base64 serialization also adds substantial peak-memory amplification beyond the raw file size.
This means the previous request to either reject the unsupported old-worker path explicitly or make it enforceably bounded has not been addressed. A receiver-side frame limit cannot provide a sender-side memory bound.
Please avoid invoking the whole-file RPC when the chunk method is unavailable. If remote log storage can provide the file, stream it from there; otherwise return an explicit “worker upgrade required for large log download” error. If compatibility for small logs must be preserved, it needs a mechanism that can establish a safe size before requesting the full payload—calling the legacy RPC first is inherently unbounded.
Please also add a rolling-upgrade regression test that verifies an old worker is never asked for the whole-file payload on the large-log path.
889a536 to
29aabe8
Compare
|
Thanks for the review. The download path no longer falls back to the legacy whole-file RPC, which prevents OOM on the worker. Changes
The RPC interface and the worker/master-side When the chunk RPC is not available, the first-chunk failure now falls back directly to If remote log storage is also unavailable, the download fails instead of trying the legacy RPC again. The error message distinguishes between:
The remote-storage error is preserved as the cause. I also added The test verifies two cases:
So the old worker is no longer asked to construct the whole log payload during the rolling-upgrade download path. I also removed the small-log compatibility fallback. There is currently no way to determine the log size safely before invoking the legacy RPC, so keeping that fallback would still leave the sender-side allocation unbounded. Log viewing is unaffected because it continues to use the line-bounded API. The changes are covered by the new integration test and the Verification40 unit/integration tests green on dolphinscheduler-api (JDK 8) plus |
There was a problem hiding this comment.
-
** Keep the remote log file stable throughout streaming**
[Code evidence]
RemoteLogClient.java:99releases the lock before streaming the file to the HTTP response. A concurrent download or log-view request can then truncate and rewrite the same file through the S3/ABS handler. The active download can hit premature EOF and fail even though the archived log is valid. A slow client extends this race window across the entire transfer.Please use a separate temporary file per download, or download to a temporary file and atomically replace the cache, so existing readers retain a stable file.
-
** Wait for async completion before asserting the response body**
[Code evidence]
LoggerControllerStreamingTest.java:125–131asserts the response body immediately after the initialMockMvc.perform().StreamingResponseBodywrites on an asynchronous thread, and the initial request does not wait for that write to finish. The assertion can therefore observe an empty response and fail intermittently; it also does not verify the final asynchronous dispatch.Please assert
asyncStarted(), capture theMvcResult, and check the final response throughasyncDispatch(mvcResult).
Replace whole-file log download with chunked streaming, and keep the legacy whole-file worker RPC unreachable from the download path: - Add ILogService#getTaskInstanceLogFileChunk RPC to read [offset, offset+length) ranges from the worker, clamped to 8 MB per chunk. - API streams chunks via StreamingResponseBody; auth is checked synchronously before the HTTP response is committed so @ApiException still returns JSON errors. - The download path never invokes the legacy whole-file worker RPC (getTaskInstanceWholeLogFileBytes): that RPC reads the entire file into the worker's heap (ByteArrayOutputStream + toByteArray + JSON/base64) before any frame exists to reject, so a large log can OOM the worker — a receiver-side maxFrameSize cannot bound the sender. On a first-chunk failure the only fallback is remote log storage, streamed in bounded chunks; the fallback runs at most once per call and its failures propagate directly, so it cannot re-enter or double-execute. Mid-stream failure throws IOException to avoid a corrupted download. - First-chunk failure handling is guided by how the RPC failed: the worker answered but could not dispatch the method (MethodInvocationException — an old worker without the chunk method, e.g. during a rolling upgrade) fails with an explicit "Worker upgrade required for large log download: chunked log RPC is not available on worker <host> and remote log storage also failed"; the worker never answered (connect refused / timeout — down or unreachable) reports a reachability problem instead of blaming the worker version; a structured non-SUCCESS response from a worker that implements the chunk RPC is reported as-is with its response code. Small-log compatibility with old workers is intentionally dropped: no mechanism can establish a byte-size bound on an old worker before the payload is built. Log viewing (line-bounded) is unaffected; the wire surface is untouched, so an old API server keeps working against a new worker/master. - The remote archive is stable throughout streaming: the remote log handlers rewrite the cache in place, so streamWholeLog snapshots the archive into a private temp file (<archive>.download-<uuid>) inside the striped lock and streams the snapshot outside it — a concurrent download/view re-downloading (and truncating) the cache can no longer cause a premature EOF for an active transfer. Snapshot deletion is guaranteed on every failure path (nested finally, up to Error); orphaned snapshots from JVM death mid-transfer are swept at startup (older than 1h, so a shared-disk instance's in-flight transfer is never touched). streamBounded's short-read guard stays as defense in depth; its message no longer blames a concurrent download (a private snapshot cannot be replaced). - Dead code removed: the byte[] download APIs (LoggerService#getLogBytes, LogClientDelegate#getWholeLogBytes, LocalLogClient#getWholeLog, RemoteLogClient#getWholeLog) and LogUtils#getFileContentBytesFromRemote in dolphinscheduler-common (zero callers after the above; the same whole-file-into-a-byte[] reader shape this PR eliminates). getFileContentBytesFromLocal stays for the worker-side legacy RPC. Tests: worker chunk RPC (range/EOF/not-found/truncated/clamp); LogClientDelegate (chunk loop, remote fallbacks, first-chunk failure distinguishing answered-but-cannot-dispatch vs unreachable vs structured non-SUCCESS, empty log terminal, mid-stream, rotation, node gone); RollingUpgradeLogStreamingIntegrationTest (real Netty wire, old-worker proxy: chunk RPC fails, whole-file RPC works and is invocation-counted — asserts the whole-file payload is never requested on the large-log path, with and without remote storage); ResponseFuture (drain isolation, set-once cause, identity removal, fail/cancel); NettyClientHandler (shared-channel concurrent drain regression, deserialize failure, timeout/interrupt leak guards); TransporterDecoder (per-field and combined frame limits); RemoteLogClient (bounded stream, empty-vs-missing, concurrent cache rewrite mid-transfer — verified to fail on the pre-snapshot code, failed transfer still deletes its snapshot, orphan sweep age gate); real-RPC integration tests (multi-chunk download + deterministic rotation); controller MockMvc (auth failure JSON; success asserts asyncStarted and the final asyncDispatch response — the body is written on the async thread, the old assertions raced it). Verified end-to-end in standalone (embedded Jetty + real Netty RPC): a 1 GB log downloads completely (byte-identical md5) with stable heap in a 1 GB JVM hosting api+master+worker together. Co-Authored-By: Claude <noreply@anthropic.com>
fce8813 to
9c53bc6
Compare
|
Thanks for the review. Both findings are addressed. 1. Keep the remote log file stable during streamingImplemented the first option:
Added a regression test, 2. Wait for async completion before asserting the response bodyImplemented as suggested. The success test now verifies |
Was this PR generated or assisted by AI?
YES. Implementation and tests drafted with assistance from Claude (Anthropic); reviewed by human.
Purpose of the pull request
getFileContentBytesFromLocalread entire files into memory with no size limit. Downloading a large task log caused OOM on the worker.This PR caps the read at 47 MB and returns a clear error for oversized logs.
Why 47 MB, not 64 MB? The
byte[]is JSON-serialized as base64 (~1.33× expansion) before RPC transmission. 47 MB raw → ~63 MB JSON body, staying under the 64 MBmaxFrameSizeinTransporterDecoder. 64 MB raw would produce ~86 MB body and be rejected byTooLongFrameException.close #18459
Brief change log
LogUtils: addMAX_LOG_DOWNLOAD_SIZE = 47 MB;getFileContentBytesFromLocalstops reading once the limit is reached.LogServiceImpl: checks file size before reading; returnsERRORwith a clear message for oversized logs instead of silently truncating.Verify this pull request
This change added tests and can be verified as follows:
LogServiceImplTest: a 48 MB file returnsERRORwith message containing "exceeds maximum download size"../mvnw spotless:checkpasses.Pull Request Notice
Pull Request Notice