Conversation
…ing event log replay
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for fixing these, @viirya. I reviewed 1df219a.
The new byte accounting and the physical line indices look correct: a faithful port of the new fetchLine was cross-checked against a reference implementation (split on \n, drop one trailing \r, keep a line iff its UTF-8 content is <= the limit, keep the physical index) on more than a million randomized inputs (1-4 byte code points, LF/CRLF/unterminated endings, limits around the boundary, Int.MaxValue) with zero mismatches.
Summary of the inline comments, most important first:
- Memory bound at the default (
ReplayListenerBus.scalaL136): the in-loop buffer bound changed from ~maxLineLength / 2chars tomaxLineLength + 1chars while the default stays512m. With the default SHS heap (-Xmx1g), an over-limit line that the current code skips now fails withOutOfMemoryErrorbefore it is rejected. Please consider256mas the default, a separate heap-oriented cap, or documenting the heap multiple. - Test coverage: nothing observes the materialization cap anymore (L139); the U+1F600 case is serialized by Jackson as the ASCII escape
\ud83d\ude00, so the surrogate branches are never exercised and no multi-byte character straddles the limit (ReplayListenerSuite.scalaL127); plus minor clue/redundancy nits (L146). - Diagnostics (L176):
JsonParseExceptionis anIOException, so the most common mid-file corruption is rethrown without the new physicalMalformed line #N. - Pre-existing gaps outside this diff that share the semantics this PR defines (fine as follow-ups): the end-event reparse in
FsHistoryProvider(see L173) and the compaction rewrite inEventFilter(see L78) still read lines with an unboundedSource.getLines(); the latter can permanently drop events of live entities after compaction. - Performance (L142, pre-existing since SPARK-59407 and kept here): the per-char locked
BufferedReader.read()makes line splitting ~14x slower than the previousSource.getLines(). - Nits: the "disables the limit" wording (
History.scalaL173), redundantoverLong ||(L148), a simpler width expression (L124), and the PR description says three new tests fail on the original implementation while four do (L112).
Since this changes the unit of a config that is already on master, branch-4.x, branch-4.3, branch-4.2, branch-4.1, branch-4.0 and branch-3.5 (all unreleased), it should land on all of those branches so that the same key has the same meaning everywhere.
| 3 | ||
| }) | ||
| // Allow one extra byte until we know whether the line ends in CRLF. | ||
| if (byteLength <= maxLineLength.toLong + 1) { |
There was a problem hiding this comment.
This changes the per-line memory bound, not only the unit. Previously the builder stopped growing at maxLineLength / 2 chars; now it can hold maxLineLength + 1 chars, and a single char >= U+0100 inflates it to UTF-16 (2 bytes per char). With the unchanged default of 512m:
- Default SHS heap (
-Xmx1g), one line of U+0101 followed by 600,000,000 ASCII chars: the current code stops at capacity 301,989,886 chars and skips the line. With this PR the builder grows to capacity 603,979,774 chars (a ~1.2 GB array) and fails withOutOfMemoryError: Java heap spacebefore the line is rejected.case e: Exceptiondoes not catch it. - For configured limits in [576m, 1152m), a Latin-1 builder that has grown past 603,979,774 chars has capacity 1,207,959,550.
AbstractStringBuilder.inflate()sizes the UTF-16 array from the capacity, so the first char >= U+0100 fails withUTF16 String size is 1207959550, should be less than 1073741823regardless of the heap size. - On JDK 8 (branch-3.5) there are no compact strings, so even pure-ASCII lines need about twice the heap of the current code.
The docs still say the setting bounds "the memory replay can use". Could we default to 256m (keeps the current bound for ASCII lines), keep a separate heap-oriented cap for the buffer, or document that heap usage can reach about 2-4x the limit?
There was a problem hiding this comment.
Confirmed. I reproduced the OOM in a separate -Xmx1g JVM using a lazy stream containing U+0101 followed by 600,000,000 ASCII characters. The 512 MiB limit fails during builder growth; the revised implementation with a 256 MiB limit successfully drains and skips the same input.
Changed both defaults to 256m in ab7c85d to keep the maximum buffered character count close to the previous effective bound. The configuration docs now describe the additional heap cost of buffer growth, UTF-16 storage and JSON parsing. This does not guarantee that every accepted 256 MiB event fits in a 1 GiB heap.
| if (byteLength <= maxLineLength.toLong + 1) { | ||
| sb.append(c.toChar) | ||
| } else { | ||
| overLong = true |
There was a problem hiding this comment.
Test coverage: because L148 now re-decides from the full byteLength, this cap is no longer observable by any test. A mutant that always appends (never sets overLong), or one that uses 2 * maxLineLength + 1, produces the same lines, indices, warnings and exceptions for every input in the suite; only the buffered size of skipped lines changes (e.g. 8193 -> 9216 chars in "Over-long event log lines are skipped instead of materialized"). Before this PR, removing the cap made that test fail. A test that observes the buffer bound would guard SPARK-59407's main guarantee, e.g. by extracting the reader into a helper and asserting the max buffered length, or by streaming a huge line from a lazy InputStream with a small limit.
There was a problem hiding this comment.
Addressed in ab7c85d. The reader now uses a BoundedLineBuffer, and the new test checks its retained character count after every append while feeding content well beyond the limit. It covers zero and small limits and allows only the one extra character reserved for a possible trailing CR.
I also tested mutations that remove the append cap or double it. Both cause this test to fail, even though the oversized line would still be rejected at the end.
| overLong = true | ||
| } | ||
| } | ||
| c = reader.read() |
There was a problem hiding this comment.
Performance (pre-existing since SPARK-59407, kept by this rewrite): BufferedReader.read() takes a lock on every call (on JDK 21 an InternalLock/ReentrantLock, whose CAS survives C2 inlining). On ~100 MB of TaskEnd-like ASCII lines this loop splits at ~73 MB/s (13.6 s/GB) vs ~1,039 MB/s (0.96 s/GB) for the previous Source.getLines(); together with Jackson readTree it is 60.9 MB/s vs 293 MB/s. That is roughly +12.6 s of CPU per GB for UI rebuilds and compaction, and ~8 s to drain a 600 MiB line. The UTF-8 width branching added here is not measurable. Keeping the same per-char logic but reading into a local char[] via InputStreamReader.read(char[], ...) is ~5x faster, and bulk-scanning the buffer for \n with slice appends is ~15x faster. Maybe a follow-up?
There was a problem hiding this comment.
Agreed that this deserves a separate follow-up. We will address block-oriented reading and validate it with an SHS replay benchmark there. This PR keeps the existing read strategy and focuses on the byte-limit and diagnostic fixes; I have not independently rerun the throughput measurements quoted here.
| sb.setLength(sb.length() - 1) | ||
| byteLength -= 1 | ||
| } | ||
| if (overLong || byteLength > maxLineLength) { |
There was a problem hiding this comment.
nit: overLong || is redundant. overLong is only set when byteLength >= maxLineLength + 2, byteLength no longer changes afterwards, and the CR strip subtracts at most 1, so byteLength > maxLineLength already holds (also for Int.MaxValue, 0 and negative limits; checked exhaustively for small limits plus random inputs). Either use if (byteLength > maxLineLength), or drop the flag and keep only if (byteLength <= maxLineLength + 1L) sb.append(c.toChar) in the loop, which removes one nesting level.
There was a problem hiding this comment.
Removed overLong in ab7c85d. The bounded accumulator stops updating once the byte count exceeds the buffer allowance, and the final content-byte check determines whether to skip the line. This also avoids continuing to count the rest of an arbitrarily long discarded line.
| if (!overLong) { | ||
| // The strict decoder guarantees valid surrogate pairs. Count their four UTF-8 | ||
| // bytes on the high surrogate and zero on the low surrogate. | ||
| byteLength += (if (c < 0x80) { |
There was a problem hiding this comment.
nit: since the strict decoder only emits complete surrogate pairs, counting 2 bytes per surrogate half gives identical results (the running total matches at every code point boundary), so this could be if (c < 0x80) 1 else if (c < 0x800 || Character.isSurrogate(c.toChar)) 2 else 3, and the comment above becomes simpler.
There was a problem hiding this comment.
Applied in ab7c85d: each surrogate half contributes two bytes, with a comment explaining the strict decoder's valid-pair guarantee. The new raw UTF-8 tests exercise supplementary code points and limits inside those code points. A mutation counting three bytes per half fails those tests.
| replayEntries(lines.zipWithIndex, sourceName, maybeTruncated, eventsFilter) | ||
| } | ||
|
|
||
| private def replayEntries( |
There was a problem hiding this comment.
The physical line number only surfaces for non-IOException failures. JsonParseException (and JsonMappingExceptions such as MismatchedInputException) extend IOException, so for a malformed JSON line in the middle of a file, throw jpe (L219) is caught by case ioe: IOException => throw ioe (L233) before the Malformed line #N log (L237). For example, [end, 2048 x's, 2048 x's, {bad, end] with a 1024 limit throws without any line number, and FsHistoryProvider only logs Jackson's line: 1, column: N, which is relative to the single-line string. {} in the new test goes through the NPE path. The case ordering is old (SPARK-2261), but logging sourceName/lineNumber before rethrowing would make this fix effective for the most common corruption.
There was a problem hiding this comment.
Confirmed and fixed in ab7c85d. JsonProcessingException now logs the source and physical line number before the same exception is rethrown. Tests cover both malformed JSON and a mapping error on physical line 4 after two skipped lines.
Truncated-last-line handling and ordinary stream IOException propagation are preserved. A separate test checks that a stream read failure does not receive a stale JSON parse-line diagnostic. Removing the new JSON diagnostic makes its regression test fail.
| .doc("Maximum UTF-8 byte length of a single event log line during replay, excluding " + | ||
| "the line ending. Longer lines are skipped with a warning, bounding the " + | ||
| "memory replay can use when an event log is corrupt or unexpectedly large. Setting " + | ||
| "this to 0 or a negative value disables the limit. " + |
There was a problem hiding this comment.
With the new Long arithmetic, "disables the limit" is no longer strictly true. ReplayListenerBus.maxLineLength maps <= 0 (and values above Int.MaxValue) to Int.MaxValue, and lines with more than 2,147,483,647 UTF-8 bytes are now skipped with a warning (e.g. 2^30 x U+00E9 is 2,147,483,648 bytes but only a 1 GiB Latin-1 builder), whereas the previous sb.length() * 2 < Int.MaxValue overflowed and never skipped. Spark 3.4+ writers cannot produce such a line (one ByteArrayOutputStream per event), so this is mostly theoretical, but skipping the check when the limit is disabled, or documenting the 2^31-1 byte cap, would keep the doc accurate.
There was a problem hiding this comment.
Corrected the configuration docs and helper comment in ab7c85d. They now state that non-positive values and values above Int.MaxValue use the maximum supported limit of 2,147,483,647 bytes. The existing configuration normalization is unchanged. Added tests for zero, negative, exact-maximum and larger values.
| test("Replay line limit handles UTF-8 boundaries and line endings") { | ||
| val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) | ||
| // scalastyle:off nonascii | ||
| val names = Seq("x", "\u00e9", "\u4e2d", "\ud83d\ude00") |
There was a problem hiding this comment.
The U+1F600 case does not exercise the 4-byte path. JsonProtocol serializes through Jackson's UTF8JsonGenerator, which writes supplementary characters as the ASCII escape \ud83d\ude00 (COMBINE_UNICODE_SURROGATES_IN_UTF8 is off by default in 2.22.2). So that JSON line is pure ASCII, the surrogate branches (ReplayListenerBus.scala L128-L131) are never taken, and mutants that count 3+3, 3+0 or 0+0 bytes for a pair all pass. Also, the boundary always falls on the trailing ASCII ","Timestamp":125,"User":"user"}, so no multi-byte character straddles the limit, and since the name starts at byte offset 110, the 8192-byte decoder boundary never splits a U+00E9/U+4E2D sequence (8192 - 110 = 8082 is divisible by 2 and 3). Could we build raw UTF-8 input directly (e.g. lines whose last code point is 2, 3 or 4 bytes) and put the limit at -1/0/+1/+2 bytes around it?
There was a problem hiding this comment.
Confirmed: Jackson escapes U+1F600 in the original test. Replaced that grid in ab7c85d with raw UTF-8 input observed through the replay filter before JSON parsing.
The final code point is 1–4 bytes wide and starts at offsets 8, 8190, 8191 or 8192. Limits range from three bytes below to two bytes above the content length, covering positions inside multibyte code points and sequences split across the reader's byte buffer. LF, CRLF and unterminated final lines are included. A mutant counting supplementary characters as 3+3 bytes fails the new grid.
| .getBytes(StandardCharsets.UTF_8)) | ||
| assert(bus.replay(input, "utf8")) | ||
| val expected = if (excess == 0) Seq(end, json) else Seq(end) | ||
| assert(listener.loggedEvents.toSeq == expected, s"name=$name ending=$ending excess=$excess") |
There was a problem hiding this comment.
nit: (1) The clue omits repetitions and prints ending raw, so \n and \r\n failures look identical, e.g. s"name=$name repetitions=$repetitions ending=${ending.map(_.toInt).mkString(",")} excess=$excess"; the assert(bus.replay(input, "utf8")) above has no clue at all. (2) repetitions 100 vs 4097 gives the same outcome for every (name, ending, excess) combination, so it only doubles the grid (see the decoder-boundary note on L127). (3) "Replay line limit uses UTF-8 bytes" is subsumed by the name = "x", excess = 0 iterations here.
There was a problem hiding this comment.
Updated in ab7c85d. Both replay success and observed-line assertions are wrapped in a clue containing the code point, padding, numeric line-ending bytes and limit delta. The old repetitions grid is replaced by deliberate boundary offsets.
I retained the short ASCII application-event test because it now also exercises the 8k configuration and successful JSON event delivery. The raw UTF-8 grid deliberately bypasses JSON parsing so serialization cannot hide the byte sequences being tested.
| assert(eventMonster.loggedEvents(1) === JsonProtocol.sparkEventToJsonString(applicationEnd)) | ||
| } | ||
|
|
||
| test("Replay line limit uses UTF-8 bytes") { |
There was a problem hiding this comment.
minor, about the PR description: it says "Three new regression tests fail against the original implementation", but four of the five new tests fail on the current code (this one, the boundary grid, the physical-line-number test and the truncated/filtered test); only "Replay still rejects malformed UTF-8 in skipped lines" passes there. Since the description becomes the commit message, it may be worth correcting.
There was a problem hiding this comment.
Confirmed by rerunning all five tests from the initial PR revision against the original implementation: four fail, and the malformed-UTF-8 test passes. The earlier count of three came from a run before the truncated/filtered test was added. I removed the stale count from the PR description.
The revised suite has 15 passing tests, and both production and test scalastyle checks pass.
HyukjinKwon
left a comment
There was a problem hiding this comment.
Reviewed ab7c85d. The change looks correct to me, and the points from the earlier review that were marked addressed hold up on this revision:
BoundedLineBufferbyte accounting is right for UTF-8 (1/2/3 bytes, and 2 per surrogate half = 4 per valid pair from the strict decoder). Retention is capped atmaxLineLength + 1bytes' worth of chars, and a trailing CR is only discounted when it is the last retained char. So an over-long line with a retained interior CR is still rejected, and a code point straddling the bound is dropped and rejected.- Physical line numbers: the index is assigned before draining, so skipped lines still consume a number. The
Iterator[String]overload keepszipWithIndex, and the filter is applied after indexing on both paths, so the truncated-last-line detection is unchanged. - The new
JsonProcessingExceptioncase only adds context. Previously these fell into theIOExceptionrethrow with no source/line.lineNumberis assigned before parsing, and reader IO errors (includingMalformedInputException) do not reach this case, which the IO test covers. - Default 256m UTF-8 bytes keeps the ASCII character capacity equal to the old
512m / 2cap. All three production callers (FsHistoryProviderx2,EventLogFileCompactor) derive the limit fromReplayListenerBus.maxLineLength(conf). Since the config was only added in SPARK-59407 and has not shipped in a release yet, redefining its unit and default here is not a compatibility concern. - The tests discriminate the fix: the 8k ASCII test fails under the old
chars * 2rule, and the raw-UTF-8 grid crosses the decoder's 8192 buffer boundary for each encoded width and line ending.
The deferred items (tail reparse via the unbounded reader, compaction reader consistency, per-char BufferedReader.read() locking) make sense as the planned follow-ups.
One small readability nit inline.
| } | ||
| if (overLong) { | ||
| val line = buffer.result() | ||
| if (line.isEmpty) { |
There was a problem hiding this comment.
nit: line is an Option[String] here, so line.isEmpty reads like an empty-string check at first glance (and an empty line is in fact accepted as Some("")). Something like buffer.result() match { case Some(line) => (line, index); case None => ... }, or naming it maybeLine, would make the skip branch clearer.
There was a problem hiding this comment.
Renamed it to maybeLine in aaee04e to make the Option[String] semantics clear. Thanks!
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for the update, @viirya. I reviewed aaee04e.
The points from the previous round that were marked as addressed hold up. I re-checked the new BoundedLineBuffer and the physical line indices with a verbatim copy of the reader: it matches a reference implementation on 900K randomized inputs (1-4 byte code points, CR/CRLF/unterminated endings, malformed UTF-8, streams split at arbitrary byte positions, limits around the boundaries and 0/negative/Int.MaxValue) with zero mismatches, and mutants of the surrogate width, the CR discount and the buffer allowance are all caught.
Summary of the inline comments, most important first:
- OOM instead of skip for large limits (
ReplayListenerBus.scalaL256): retaininglimit + 1chars can exceedStringBuilder's implementation limits (Latin-1 to UTF-16 inflation from 576m, the UTF-16 maximum length from about 1g), soOutOfMemoryErroris thrown regardless of the heap size where the current code skips the line, and even for some lines within the limit. On the listing path the error is swallowed silently. Clamping the effective limit (e.g.512m) and documenting it as the maximum would fix it. - Backports (
History.scalaL181): this needs to land on branch-4.3 before the next 4.3.0 RC (v4.3.0-rc1 (2026-08-31) does not contain SPARK-59407, butbranch-4.3does), and branch-3.5 needs a separate PR because the new structured-logging block merges cleanly there but does not compile. - Duplicate stack traces (L210): the new log passes
jpealthough every caller already logs or retries the same exception, so UI loads now add 2-4 ERROR stack traces per request. Logging only the location would be enough. - Test coverage (
ReplayListenerSuite.scalaL259, pre-existing): the end-to-end replay tests comparee1withe1and ignore the result, so the "existing compressed/end-to-end replay tests" in the description don't check the replayed events. - Tail-reparse follow-up (L148): the new ERROR prints skip-relative line numbers on that path, and (pre-existing) a skip that lands inside a multi-byte sequence drops a completed app from the listing, which reusing
boundedLineswould not fix by itself. - Truncation inside a code point (L83, pre-existing): an in-progress log cut inside a multi-byte sequence fails with
MalformedInputExceptioninstead of being tolerated as truncated. - Nit (L239): derive
DEFAULT_MAX_LINE_LENGTHfrom the config default and describe the current bound in its comment.
Items 4-6 are pre-existing and fine as follow-ups.
| private[scheduler] class BoundedLineBuffer(maxLineLength: Int) { | ||
| private val buffer = new java.lang.StringBuilder() | ||
| private var byteLength = 0L | ||
| private val bufferLimit = maxLineLength.toLong + 1 |
There was a problem hiding this comment.
This retains up to maxLineLength + 1 chars, but nothing caps that against the StringBuilder implementation limits, so larger configured limits now fail with OutOfMemoryError, regardless of -Xmx, where the current code skips the line:
- Latin-1 to UTF-16 inflation (limits >= 603,979,775 bytes, i.e. from 576m). Once more than 603,979,774 Latin-1 chars are retained, the builder's capacity is 1,207,959,550, and the first char >= U+0100 makes
inflate()fail withUTF16 String size is 1207959550, should be less than 1073741823. This also hits lines within the limit: with600m, 620,000,000xfollowed by U+0100 (620,000,002 bytes) throws, while the current code skips it. The current code needs 1152m or more for this. It is the [576m, 1152m) case from the first round, which the new default avoids only for the default. - UTF-16 maximum length (limits >= 1,073,741,823 bytes, about 1g). A UTF-16 builder cannot exceed 1,073,741,822 chars, so e.g.
€followed by about 1.18 billionxat1gthrowsRequested array size exceeds VM limit, while the current code skips it. <= 0and values aboveInt.MaxValuefail the same way as the current code, but the new doc now calls 2147483647 bytes the "maximum supported limit".
On the listing path, the Error passes mergeApplicationListing's case e: Exception and is swallowed by the Future of pool.submit in submitLogProcessTask, so the app is silently missing from the listing. 256m and an explicit 512m are safe.
Could we clamp the effective limit in maxLineLength(conf) (and in the constructor if needed), e.g. to 512m (any value up to 603,979,773 keeps the capacity at or below 603,979,774), and document that as the maximum? The clamp needs to apply to the limit itself, not only to bufferLimit, so that lines within the limit are not silently truncated. I reproduced these with a verbatim copy of the reader on JDK 21 (up to -Xmx8g).
There was a problem hiding this comment.
I reproduced the 600m / 620,000,000 ASCII characters + U+0100 case with the actual reader on JDK 21 and -Xmx3g. The effective content limit is now capped at 512m in both configuration resolution and the constructor, and the same input is successfully skipped. Tests cover normalization and the docs distinguish this JVM limit from heap requirements.
| .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) | ||
| .bytesConf(ByteUnit.BYTE) | ||
| .createWithDefaultString("512m") | ||
| .createWithDefaultString("256m") |
There was a problem hiding this comment.
A reminder about the backports from the first round: this changes the unit and the default of a key that branch-4.x, branch-4.3, branch-4.2, branch-4.1, branch-4.0 and branch-3.5 already have (all still 512m with the chars * 2 check, and not in any release yet).
- v4.3.0-rc1 (2026-08-31) does not contain SPARK-59407, but
branch-4.3does, so this needs to land onbranch-4.3before the next 4.3.0 RC. Otherwise 4.3.0 ships the old semantics (e.g.8kmeans about 4096 chars and<= 0means no limit), and changing it afterwards becomes a behavior change of a released config. - The diff cherry-picks cleanly to branch-4.x and branch-4.3, and branch-4.2/4.1/4.0 conflict only in the scaladoc (
byes/bytes). On branch-3.5, the newcase jpe: JsonProcessingExceptionblock withlog"..."/MDCmerges without a conflict but does not compile, so it needs a separate PR like [SPARK-59407][UI][3.5] Improve compressed line handling in Spark History server #59058.
There was a problem hiding this comment.
Confirmed: v4.3.0-rc1 does not contain SPARK-59407, while branch-4.3 does. This fix needs to reach branch-4.3 before the next RC, along with the other affected branches. Branch-3.5 needs a separate backport adapting the logging API.
| case _: EOFException if maybeTruncated => false | ||
| case jpe: JsonProcessingException => | ||
| logError(log"Exception parsing Spark event log: ${MDC(PATH, sourceName)} " + | ||
| log"at line ${MDC(LINE_NUM, lineNumber)}", jpe) |
There was a problem hiding this comment.
Thanks for adding the location. Passing jpe here duplicates the stack trace, though, since every caller already logs or retries the same IOException:
- UI loads:
createHybridStore,createDiskStoreandcreateInMemoryStore(FsHistoryProvider L1598, L1672, L1701) retry once onIOException, so one malformed line now adds 2 ERROR stack traces per request (4 with the hybrid store), where the current code logs none at ERROR (only the INFO "trying again"). Failed loads are not cached, so this repeats on every refresh. - Listing (L1032) and compaction (L1224) log the same exception again at ERROR.
Could we drop the throwable and keep only the source and line (optionally with jpe.getOriginalMessage)? The new test only checks the message, so it would still pass.
There was a problem hiding this comment.
Removed the throwable from the location diagnostic and kept the same exception rethrow. The tests now assert a single location diagnostic with getThrown == null for both parsing and mapping errors.
| } | ||
|
|
||
| /** | ||
| * Test replaying compressed spark history file that internally throws an EOFException. To |
There was a problem hiding this comment.
About the "existing compressed/end-to-end replay tests" in the PR description: they don't really check the replayed events (pre-existing, not from this PR).
testApplicationReplay, used by "End-to-end replay" and "End-to-end replay with compression", compares each event with itself,JsonProtocolSuite.assertEquals(e1, e1)(L408). SPARK-28770 changed it from(e1, e2), which looks like a typo. It alsozips without a size check and ignores the result ofreplay(L396), so it passes even if nothing is replayed.- This test only checks the event count (L300) and also ignores the result of
replay(L298).
So nothing checks the content of a multi-buffer LZ4/Snappy/ZSTD log read through the new reader. Could we assert the result, compare originalEvents.size with replayedEvents.size and use assertEquals(e1, e2) (and pick the current app's log instead of the latest file in the shared directory), here or in a small test-only follow-up? Otherwise, please adjust the claim in the description.
There was a problem hiding this comment.
Confirmed the pre-existing assertion gaps. I have narrowed the testing description; we will fix the replay result, event count/content assertions and log-file selection in a separate test follow-up.
| sourceName: String, | ||
| maybeTruncated: Boolean, | ||
| eventsFilter: ReplayEventsFilter): Boolean = { | ||
| replayEntries(lines.zipWithIndex, sourceName, maybeTruncated, eventsFilter) |
There was a problem hiding this comment.
Two notes for the tail-reparse follow-up (FsHistoryProvider L1104-L1127), which calls this overload:
- There,
linesstarts after the byte skip and the dropped partial line, so the new ERROR at L208-L211 prints the full path with a skip-relative number. For example, a truncatedSparkListenerApplicationEndline at physical line 12,004 of a completed app is reported asat line 3959. Before this PR, JSON errors on this path had no bus-level line number. Please cover it in the follow-up, together with the existingMalformed line #Nand the truncation WARN, which are also skip-relative there, or label the source on that path, e.g.s"$path (lines counted after skipping $target bytes)". - Pre-existing: if
targetlands inside a multi-byte UTF-8 sequence (e.g. CJK SQL text or job descriptions, which Jackson writes as raw UTF-8),source.next()(L1124) throwsMalformedInputExceptionbeforereplayis called.mergeApplicationListingonly logs it, and since the placeholderLogInfowas already written with the same file size,shouldReloadLognever retries it, so the completed app stays missing from the listing (andtargetis deterministic, so a restart fails the same way). Replaying throughboundedLinesafter the skip, as planned, would throw the same way because of its strict decoder, so the follow-up should drop the bytes up to the first\nbefore decoding. I reproduced it at the unit level, not yet on a real SHS.
There was a problem hiding this comment.
Confirmed the skip-relative numbering and reproduced source.next() failing when the offset lands inside a UTF-8 sequence. We will cover both in the tail-reparse follow-up, including discarding the partial line as raw bytes before decoding and correcting the diagnostic location. The SHS listing/retry behavior still needs an integration regression test.
| sourceName: String): Iterator[(String, Int)] = { | ||
| // Fail on malformed input like Source.getLines() does instead of replacing it. | ||
| val decoder = StandardCharsets.UTF_8.newDecoder() | ||
| .onMalformedInput(CodingErrorAction.REPORT) |
There was a problem hiding this comment.
Pre-existing and preserved by this PR (Source.getLines() with Codec.UTF8 also fails on malformed input), but related to the truncation handling: if the last, unterminated line of an in-progress or crashed app's log ends inside a multi-byte UTF-8 sequence, this decoder throws MalformedInputException from reader.read() at EOF. That happens in lineEntries.hasNext, outside the inner try, so it never reaches the maybeTruncated tolerance for JsonParseException; it is rethrown as an IOException, and the UI rebuild fails after its retry (permanently for a crashed app). The same log cut at an ASCII boundary replays fine.
Compressed in-progress logs expose lengths at compression block boundaries (128 KiB blocks for zstd, 32 KiB for lz4) rather than at code point boundaries, so with CJK-heavy content this happened in 3/80 (zstd) and 5/80 (lz4) simulated cuts; uncompressed local files were always cut at code point boundaries in my runs. A follow-up could treat a MalformedInputException raised after the underlying stream hit EOF as a truncated last line when maybeTruncated is set.
There was a problem hiding this comment.
Reproduced: an ASCII-truncated final JSON line is tolerated with maybeTruncated, while a partial UTF-8 code point at EOF throws MalformedInputException. We will handle this separately, including compressed-log regression coverage.
| * Keep the maximum buffered character count close to the original 512 MiB / 2 cap. | ||
| */ | ||
| val DEFAULT_MAX_LINE_LENGTH: Int = 512 * 1024 * 1024 | ||
| val DEFAULT_MAX_LINE_LENGTH: Int = 256 * 1024 * 1024 |
There was a problem hiding this comment.
nit: This duplicates the 256m default of History.EVENT_LOG_MAX_LINE_LENGTH, and this PR had to change both. Deriving it, e.g. History.EVENT_LOG_MAX_LINE_LENGTH.defaultValue.get.toInt (like ChunkedByteBuffer does with BUFFER_WRITE_CHUNK_SIZE; the API is the same down to branch-3.5), keeps them in sync. Also, "the original 512 MiB / 2 cap" refers to the unreleased code that this PR replaces, so it will be hard to interpret after the merge. Could the comment describe the current bound instead (a skipped line retains at most this many chars plus a possible trailing CR while it is drained) and leave the history to the PR description?
There was a problem hiding this comment.
The reader default now comes from History.EVENT_LOG_MAX_LINE_LENGTH.defaultValue. I also replaced the historical comment with a description of the shared default.
What changes were proposed in this pull request?
Fix the UTF-8 byte limit and physical-line diagnostics introduced by SPARK-59407 (#58700):
512mto256m, keeping the worst-case buffered character count close to the previous implementation's effective default bound.512min both configuration resolution and the reader constructor, and document the heap overhead. Derive the reader default from the History configuration.The iterator-based replay API, strict UTF-8 decoding, existing truncated-last-line handling and ordinary input-stream IOException behavior are preserved. Independent follow-ups will address the unbounded compaction-filter and tail-reparse paths, tail-relative diagnostics and UTF-8 alignment, incomplete UTF-8 at EOF, bulk-read performance, and the pre-existing end-to-end test assertions.
JIRA: https://issues.apache.org/jira/browse/SPARK-59804
Why are the changes needed?
The original reader compares
StringBuilder.length() * 2with a byte-valued setting. With an 8 KiB limit, it incorrectly skips a valid ASCII application-start event with a 6,000-character application name. UTF-8 multibyte characters and CRLF endings also produce inconsistent boundaries.Dropping oversized lines before
zipWithIndexrenumbers subsequent diagnostics. For example, a valid event on line 1, oversized lines 2 and 3, and an invalid event on line 4 report line 2. JSON processing exceptions also need an explicit diagnostic before they are rethrown to retain that physical location.Correcting the unit without adjusting the default permits roughly twice as many buffered characters. With a lazy line containing U+0101 followed by 600,000,000 ASCII characters, the UTF-8 reader with a 512 MiB limit runs out of memory in a separate 1 GiB heap JVM; the revised 256 MiB default successfully drains and skips that input. Buffer growth, UTF-16 storage and JSON parsing still require heap beyond the configured content limit.
A separate
-Xmx3gJVM also reproduces a JVM array-size failure with600mand 620,000,000 ASCII characters followed by U+0100, despite the line fitting within that configured limit. StringBuilder capacity growth followed by UTF-16 inflation fails independently of available heap. Clamping the effective content limit to512msuccessfully skips that input; it does not guarantee adequate heap for every workload.Does this PR introduce any user-facing change?
Yes. Events at or below the effective UTF-8 content limit are replayed, LF/CRLF endings do not consume the limit, and parsing diagnostics identify physical lines after skipped content. The default is reduced from
512mto256mto keep its maximum buffered character count close to the previous effective bound.Non-positive values and values above
512mnow resolve to an effective maximum of 536,870,912 bytes (512 MiB), including direct constructor arguments. This clamps the content limit itself, so oversized lines are skipped instead of replayed as truncated prefixes. Smaller positive limits are unchanged.How was this patch tested?
All 16 suite tests and both scalastyle checks pass. Coverage includes the 8 KiB ASCII configuration case, raw 1- through 4-byte UTF-8 input with length and read-buffer boundary splits, line endings, directly observed buffer bounds, configuration and constructor normalization, physical-line diagnostics without an attached throwable, JSON parsing and mapping exceptions, filtered iterators, truncation, and malformed UTF-8. The suite also runs its existing compressed/end-to-end tests, but their pre-existing assertion gaps mean they do not establish full replayed-event content equality; those test fixes remain a separate follow-up.
Targeted mutation checks confirmed that removing or doubling the buffer cap, miscounting surrogate pairs, or removing JSON line diagnostics each makes its corresponding test fail. The oversized lazy-stream scenarios were also exercised in separate
-Xmx1gand-Xmx3gJVMs as described above.Was this patch authored or co-authored using generative AI tooling?
Generated-by: OpenAI Codex (GPT-6)