Skip to content

[SPARK-60109][UI] Apply the replay line-length limit and UTF-8-safe line alignment to the History Server end-event reparse - #59327

Open
viirya wants to merge 4 commits into
apache:masterfrom
viirya:SPARK-60109
Open

viirya wants to merge 4 commits into
apache:masterfrom
viirya:SPARK-60109

Conversation

@viirya

@viirya viirya commented Oct 9, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

When building the application listing, FsHistoryProvider looks for the application end event by reopening the last event log file, skipping to len - spark.history.fs.endEventReparseChunkSize bytes and re-parsing the rest. This PR changes how that tail is read:

  • The new FsHistoryProvider.skipToLineStart positions the stream at the first line that starts at or after the skip target, working on raw bytes before any UTF-8 decoding. 0x0A never occurs inside a multi-byte UTF-8 sequence, so this is safe at any byte offset.
    • It skips to target - 1 and scans from there, so a line that starts exactly at target (possibly the end event itself) is kept instead of being discarded.
    • When skip returns 0, it reads one byte to tell EOF apart from a short skip, so it terminates when the decompressed stream is shorter than target (e.g. when the compressed length exceeds the decompressed one).
    • The partial line is scanned in bulk; bytes read past the line boundary are put back in front of the stream with a SequenceInputStream.
  • The tail is then replayed through the bounded InputStream reader instead of the Iterator[String] overload fed by Source.fromInputStream(in)(Codec.UTF8).getLines(), so spark.history.fs.eventLog.maxLineLength and the strict decoder apply on this path as well.
  • The new ReplayListenerBus.replayFromOffset takes the byte offset of the stream's first line. When it is positive, every replay diagnostic that prints a line number (JSON errors, Malformed line #N, the truncation warning and the over-long line warning) adds (lines counted from uncompressed byte offset N), with the offset logged under its own OFFSET key. The source name stays the plain file path, so the PATH / FILE_NAME fields are unchanged. The existing replay signatures are unchanged.

maybeTruncated, eventsFilter and HaltReplayException handling are unchanged, since all replay variants share replayEntries. Without a skip (target <= 0), replay starts at offset 0 and the diagnostics are as before.

Why are the changes needed?

This is the tail-reparse follow-up discussed in #59074. With the defaults (completed apps, endEventReparseChunkSize=1m), every completed application takes this path, and because the target is computed from the compressed file length while the skip runs on the decompressed stream, most of a compressed log can be read there.

  1. The Iterator[String] overload bypasses the bounded line reader (SPARK-59407 / SPARK-59804), so maxLineLength did not apply and an over-long line was fully materialized.
  2. If the skip target landed inside a multi-byte UTF-8 character (e.g. CJK SQL text or job descriptions, which Jackson writes as raw UTF-8), source.next() threw MalformedInputException before replay started. mergeApplicationListing only logs it, and since the placeholder listing entry was already written with the same file size, the log is not retried, so the completed application stayed missing from the listing. The target is deterministic, so restarting the History Server failed the same way.
  3. Line numbers in the replay diagnostics were relative to the skip point but printed next to the full file path.

Because the reparsed tail of a compressed log is most of the file, routing it through the bounded reader makes this step several times slower with the current per-character reader. #59360 (SPARK-60158) makes the bounded reader read in bulk and should be merged before this PR; the two should be backported together.

Does this PR introduce any user-facing change?

Yes, in the History Server:

  • A completed application whose end-event reparse offset lands inside a multi-byte UTF-8 character now appears in the application listing instead of being dropped.
  • spark.history.fs.eventLog.maxLineLength now also applies to the end-event reparse, so over-long lines there are skipped with a warning instead of being materialized.
  • Replay log messages on this path name the byte offset that their line numbers are relative to.
  • A line that starts exactly at the reparse offset is no longer discarded, so an end event there no longer forces a reparse of the whole log.
  • If the decompressed stream ends before the reparse offset, the scan no longer spins forever in the skip loop; the replay sees an empty input (previously, reaching EOF exactly at the offset failed with NoSuchElementException from source.next()).

How was this patch tested?

New tests in FsHistoryProviderSuite:

  • end event reparse when the skip offset lands inside a multi-byte character: the reparse offset is set to the third byte of a 3-byte character, so the line boundary scan starts on a continuation byte. Before this PR, the application is missing from the listing.
  • end event reparse keeps a line that starts exactly at the skip offset: the offset is the first byte of the end event, and a malformed line before it makes the whole-log fallback reparse fail, so the app is listed as completed only if the tail reparse finds the end event.
  • end event reparse stops at the end of a stream shorter than the skip target: a tiny lz4 log whose compressed file is larger than its content, with a 1-byte chunk; the scan must finish (it runs with a timeout).
  • end event reparse applies the line length limit and reports relative line numbers (uncompressed and lz4): an over-long line in the tail is skipped under a 1k limit, and the warning names the relative line number and the offset in the uncompressed stream.
  • end event reparse reports JSON errors with the skip offset: a malformed end event in the tail is reported with the relative line number and the offset, the path field stays the plain file name, and the application is not listed.
  • skipToLineStart aligns to the first line starting at or after the target: targets inside a line, at its \n, exactly at a line start, past EOF and non-positive; on each byte of a 3-byte character; CRLF; no trailing newline; a partial line longer than the read buffer; and a stream whose skip never advances.

New test in ReplayListenerSuite:

  • Replay from an offset reports where line numbers start: replayFromOffset with a positive offset labels the over-long line warning, Malformed line #N and the truncated-file warning, and offset 0 leaves the diagnostics unchanged.

FsHistoryProviderSuite (RocksDB backend), ReplayListenerSuite, core/scalastyle and core/Test/scalastyle pass on master.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Code (Claude Opus 5.5)

…ine alignment to the History Server end-event reparse

Co-authored-by: Claude Code

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for working on this, @viirya. Moving the tail reparse onto the bounded InputStream replay and aligning on raw bytes looks like the right direction. I left 8 inline comments; here is a summary.

Correctness

  1. The unchanged skip loop never terminates when the decompressed stream is shorter than target, because the codec streams return 0 from skip at EOF. The PR description's "empty input at end of stream" only holds when the skip stops exactly at EOF.
  2. skipPartialLine always discards a line, so a target that lands exactly at a line start throws away a complete line, possibly the end event. With raw bytes this is now cheap to avoid.
  3. UTF-8 safety is fixed only at the start of the tail. A log being written whose last line ends mid multi-byte character still fails with MalformedInputException, which maybeTruncated does not cover.

Design / efficiency

  1. Encoding the offset into sourceName puts non-path text into the PATH / FILE_NAME MDC fields of the structured logs.
  2. skipPartialLine reads one byte per call with no bound, which is costly when the partial line is very long.

Tests

  1. The expected offset is computed from String.length (UTF-16 chars) rather than UTF-8 bytes.
  2. There is no coverage for compressed logs, for the skip reaching EOF, or for skipPartialLine edge cases.
  3. The JSON-error test checks only the log message, not the resulting listing state.

log" from ${MDC(PATH, logPath)}...")
var skipped = 0L
while (skipped < target) {
skipped += in.skip(target - skipped)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[1] This loop never terminates when the stream ends before target.

This is pre-existing, but it is in the block this PR rewrites, and it contradicts the PR description ("If the skip reaches the end of the stream, the replay now sees an empty input"). target comes from the compressed getLen, while the skip runs on the decompressed stream. If the decompressed stream is shorter than target, skip returns 0 at EOF (e.g. LZ4BlockInputStream.skip returns 0 once finished, and BufferedInputStream / InputStream.skip behave the same way), so skipped stops growing and the scan thread spins forever.

Examples: a small endEventReparseChunkSize with a tiny lz4/zstd log whose framing overhead makes the compressed file larger than its content, or a log replaced by a shorter file after listing. Could we stop when skip returns 0 and confirm EOF with in.read() == -1? Then the "empty input" behavior holds in every case.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 3b54157. skipToLineStart now reads one byte when skip returns 0: -1 ends the skip at EOF, otherwise the byte counts as skipped. The replay then sees an empty input, as the description says. Covered by the skipToLineStart test with a stream whose skip always returns 0, both before EOF and with a target past EOF.

// Because skipping may leave the stream in the middle of a line, or even in the
// middle of a multi-byte UTF-8 character, discard the rest of that line before any
// decoding. Line numbers reported by the replay are then relative to this offset.
val offset = target + skipPartialLine(in)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[2] A line starting exactly at target is always discarded.

skipPartialLine always throws away bytes up to the next \n. When target lands exactly on a line start (byte target - 1 is \n), it discards a complete line, which can be the SparkListenerApplicationEnd line itself. The app then misses its end event and falls back to the slower path. This is pre-existing, but now that the code works on raw bytes it is cheap to avoid: skip target - 1 bytes, read one byte, and call skipPartialLine only if that byte is not \n.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 3b54157. It now skips to target - 1 and scans from there, so a line that starts exactly at target is kept. The new test end event reparse keeps a line that starts exactly at the skip offset puts the target on the first byte of the end event and asserts that the app is listed as completed without the "end event was not found" whole-log reparse; it fails if the skip goes to target instead.

// middle of a multi-byte UTF-8 character, discard the rest of that line before any
// decoding. Line numbers reported by the replay are then relative to this offset.
val offset = target + skipPartialLine(in)
s"${lastFile.getPath} (lines counted from uncompressed byte offset $offset)"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[4] The offset label overloads sourceName.

sourceName is logged under MDC(PATH, ...) and MDC(FILE_NAME, ...) in ReplayListenerBus. With this change, structured logs on this path get values like hdfs://.../app-1 (lines counted from uncompressed byte offset 123456) in the path field, which no longer matches the file. Tooling that groups or links logs by that key sees a different "path" for each offset. Could we pass the offset to replay separately (or log it under its own key) and keep the path field clean?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changed in 3b54157. The source name is now the plain path. The new replayFromOffset takes the start offset separately, and the diagnostics that print a line number append (lines counted from uncompressed byte offset N) with the offset under its own OFFSET key. The public replay signatures are unchanged. The JSON-error test also asserts that the path is no longer followed by the label.

}

bus.replay(source, lastFile.getPath.toString, !appCompleted, eventsFilter)
bus.replay(in, sourceName, !appCompleted, eventsFilter)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[3] UTF-8 safety is fixed only at the start of the tail.

replay(InputStream, ...) uses the strict decoder, and a MalformedInputException there is an IOException, which replayEntries rethrows regardless of maybeTruncated. With spark.history.fs.inProgressOptimization.enabled=false, an in-progress log takes this path while the writer may have flushed only the first bytes of a multi-byte character at EOF. The listing for that app then fails on this scan even though maybeTruncated = true is meant to tolerate a partial last line. Should a malformed sequence at EOF be handled like a truncated last line when maybeTruncated is set? A separate JIRA is also fine, since the same applies to the main replay path.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed, and as you note this also affects the main replay path and was the same with Source.getLines(). I filed SPARK-60114 for it, so this PR keeps the decoding behavior unchanged.

private[history] def skipPartialLine(in: InputStream): Long = {
var discarded = 0L
var b = in.read()
while (b != -1 && b != '\n') {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[5] One read() call per byte, with no bound.

in is a BufferedInputStream (or a codec stream), so each read() is a synchronized method call. If the skip lands at the start of a very long line (a large SQL plan, or a corrupted line far above maxLineLength), this makes one call per byte, and it repeats on every rescan of that log. A bulk read into a byte array that scans for \n and passes the leftover bytes back in front of the stream (e.g. new SequenceInputStream(new ByteArrayInputStream(rest), in)) would avoid that while keeping the alignment exact.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changed in 3b54157. The partial line is now scanned in 8 KiB reads, and the bytes after the \n are put back in front of the stream with SequenceInputStream, as suggested. The skipToLineStart test includes a partial line longer than the buffer.

val bytes = writeRawLog(log, lines)
// Skip into the middle of the line before the tail. The next line starts at `offset` and
// the over-long line is the first line after it.
val offset = lines.takeWhile(_ != tail.head).map(_.length + 1).sum

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[6] The expected offset is computed in UTF-16 chars, not UTF-8 bytes.

_.length + 1 counts chars, while the production offset counts UTF-8 bytes. This matches today only because every line before the tail is ASCII. If a line in endEventTestLines ever contains non-ASCII text (for example, reusing the CJK description from the first test), the expected offset would diverge. Using _.getBytes(StandardCharsets.UTF_8).length + 1 here and at L1295 would match what is being verified.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 3b54157; the tests now compute line lengths with getBytes(UTF_8).length + 1.

}
}

test("end event reparse when the skip offset lands inside a multi-byte character") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[7] Coverage gaps.

All three new tests use uncompressed logs. According to the PR description, most completed apps take the compressed path (the target comes from the compressed length and the skip runs on the decompressed stream), but none of these tests exercise it, so the "uncompressed byte offset" label is not verified there. The claimed empty-input behavior when the skip reaches EOF is also untested (see [1]), and so are the skipPartialLine edge cases: no trailing \n, \r\n, and target exactly at a line start. A small direct test of skipPartialLine plus one compressed-codec variant would cover most of this.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added in 3b54157:

  • The over-long line test now also runs with lz4, using a target derived from the compressed length. The expected offset is computed from the uncompressed bytes.
  • A direct skipToLineStart test covers targets inside a line, at its \n, exactly at a line start and past EOF; CRLF; no trailing newline; a partial line longer than the read buffer; and a stream whose skip never advances.

val conf = createTestConf().set(END_EVENT_REPARSE_CHUNK_SIZE, bytes.length - offset + 10L)
val appender = new LogAppender("end event reparse")
withLogAppender(appender) {
new FsHistoryProvider(conf).checkForLogs()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[8] Only the log message is asserted here.

This test checks the log text, but not what happens to the application listing after the JsonParseException propagates. A change that silently swallowed the malformed end event (listing the app as incomplete, or dropping it) would still pass. Could we also assert the expected listing state, e.g. via provider.getListing()?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added in 3b54157: the test now also asserts the listing. Since the log is completed (maybeTruncated = false) and the malformed line is not the last one, the JsonParseException propagates and the application is not listed, which is unchanged by this PR.

…lk alignment, offset in its own log field, more tests

Co-authored-by: Claude Code

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for addressing all the previous comments, @viirya. The new skipToLineStart / replayFromOffset design looks good to me, and deferring the EOF decoding issue to SPARK-60114 makes sense.

I left 4 minor, non-blocking inline comments:

  1. The target > 0 precondition of skipToLineStart is only documented. A non-positive target silently drops the first line.
  2. The boundary test proves the reparse found the end event only through the absence of an INFO message.
  3. Two of the four diagnostics that now append the offset label have no test.
  4. The skip-past-EOF fix is tested only with a fake stream, not with a real codec stream.

* boundary is found on raw bytes before any decoding: '\n' never occurs inside a multi-byte
* UTF-8 sequence. Bytes read past the boundary are put back in front of the stream.
*/
private[history] def skipToLineStart(in: InputStream, target: Long): (InputStream, Long) = {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[1] Minor: the target > 0 precondition is only documented.

With target <= 0, the skip loop does nothing, but the scan still discards everything up to the first \n. The helper would then silently drop line 0 and return the length of that line as the offset. The only caller guards this today, but returning (in, 0L) for target <= 0 (or adding require(target > 0)) would make the contract explicit. The first option would also let the caller drop its else (in, 0L) branch.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changed in e0f71e5. skipToLineStart now returns (in, 0L) for a non-positive target, and the caller calls it unconditionally (the else branch is gone). The skipToLineStart test covers 0 and -1.

}
// The end event is found by the reparse, not by the fallback that parses the whole log.
val messages = appender.loggingEvents.map(_.getMessage.getFormattedMessage)
assert(!messages.exists(_.contains("since end event was not found")), messages)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[2] Minor: this relies on the absence of a log message.

The list assertions above pass either way, because the whole-log fallback also lists the app as completed. So this negative match on the INFO text since end event was not found is the only thing that detects a regression. If that message is reworded or its level changes, the assertion passes vacuously. A positive signal that does not depend on unrelated wording would make the test more robust.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changed in e0f71e5. The log now has a malformed SparkListenerApplicationEnd line between the environment update and the end event. The listing pass halts before it and the tail reparse starts after it, but the whole-log fallback fails on it, so the app is listed as completed only if the tail reparse finds the end event. The log-message check is removed; the test fails if the skip goes to target instead of target - 1.

}

/** Describes where line numbers start when the replay does not start at the source start. */
private def lineOrigin(startOffset: Long): MessageWithContext = {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[3] Minor: two of the four lineOrigin call sites are untested.

The over-long line warning and the JsonProcessingException error are covered in FsHistoryProviderSuite. The Malformed line #N error and the truncated-file JsonParseException warning are not, and ReplayListenerSuite has no direct test for replayFromOffset. A small ReplayListenerSuite case that calls replayFromOffset with a positive offset on malformed and truncated input would cover both.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added Replay from an offset reports where line numbers start to ReplayListenerSuite in e0f71e5. It calls replayFromOffset with offset 42 and checks the label on the over-long line warning, Malformed line #N and the truncated-file JsonParseException warning, and that offset 0 leaves the diagnostics unlabeled. The JsonProcessingException error is covered by the FsHistoryProviderSuite test.

}
}
assert(alignStream(noSkip("abc\ndef"), 2) === (("def", 4L)))
assert(alignStream(noSkip("abc\n"), 100) === (("", 4L)))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[4] Minor: the EOF-before-target path is covered only by a fake stream.

This FilterInputStream whose skip always returns 0 checks the loop logic well. The real trigger, though, is a codec stream (lz4/zstd) whose decompressed length is shorter than target. An integration case with a tiny lz4 log and END_EVENT_REPARSE_CHUNK_SIZE = 1 (so that the compressed length exceeds the decompressed one) would pin down the end-to-end behavior that motivated the fix.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added end event reparse stops at the end of a stream shorter than the skip target in e0f71e5. It writes a one-line lz4 log whose compressed file is larger than its content (asserted in the test), and uses END_EVENT_REPARSE_CHUNK_SIZE = 1, so LZ4BlockInputStream.skip returns 0 before target - 1. The scan runs in a separate thread with a one-minute timeout, so a non-terminating skip loop fails the test instead of hanging it; it does fail that way when the EOF check is removed.

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1, LGTM (Pending CIs).

…, offset diagnostics and codec EOF tests

Co-authored-by: Claude Code

@HyukjinKwon HyukjinKwon left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review summary

The rewrite of the end-event tail reparse looks correct. I traced skipToLineStart through targets inside a line, at \n, at CRLF, exactly at a line start, past EOF, and over a partial line longer than the 8 KiB buffer. The skip() == 0 fallback correctly separates EOF from a short skip, the leftover bytes are replayed in order through SequenceInputStream, and both existing replay overloads behave as before at offset 0. The comments from the earlier review rounds are addressed at this head.

Two small follow-ups. The skipToLineStart scaladoc says in is positioned at the line start, but only the returned stream is. The multi-byte reparse test now starts its raw scan on the character's lead byte (because of the target - 1 skip), so it no longer exercises a mid-character start. I also left a question about the listing-path cost of routing the compressed-log tail through the per-char bounded reader.

Findings

2 total: 0 P0, 0 P1, 0 P2, 2 P3.

Nit (P3)

  • skipToLineStart scaladoc says in is positioned at the line start — core/src/main/scala/org/apache/spark/deploy/history/FsHistoryProvider.scala:1765 — see inline.
  • Multi-byte reparse test no longer starts the scan inside a character — core/src/test/scala/org/apache/spark/deploy/history/FsHistoryProviderSuite.scala:1258 — see inline.

Decision challenges

Listing-path cost of the per-char bounded reader on compressed logs

Routing the tail through replayFromOffset means the end-event reparse now reads lines with boundedLines, which calls BufferedReader.read() and BoundedLineBuffer.append once per character, instead of readLine via Source.getLines. For uncompressed logs the tail is only about endEventReparseChunkSize, so this doesn't matter. For compressed logs, though, target is the compressed length minus the chunk, applied to the decompressed stream, so the tail is most of the log. Since spark.eventLog.compress defaults to true, that is the common case for completed apps. The first listing pass halts after the environment update, so this read is most of the per-log listing cost.

A rough local microbenchmark of the two loops (200 MB of ASCII event-like lines, JDK 17 and 25) gave about 0.2 s for readLine versus 2.3-4.7 s for the per-char loop, roughly 10-20x on line reading. That is an approximation, not a History Server measurement. For large compressed logs, though, it could make listing noticeably slower: at startup without a persisted store, and whenever a log is new or updated. That works against what endEventReparseChunkSize is documented to do.

Is that cost acceptable here, given that full replays already pay it? Or would it be worth making boundedLines scan in bulk (read(char[]) and search for \n, as readLine does), so both paths keep readLine-like throughput while still enforcing the limit? Doing that as a separate change would also be fine.

Generated by Omnigent on Databricks.

private[spark] object FsHistoryProvider {

/**
* Positions `in` at the start of the first line that begins at or after byte `target`,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (P3): in itself isn't left at the line start. When a \n is found, in has already been read up to 8 KiB past the boundary (in.read(buffer)), and only the returned SequenceInputStream(rest, in) starts at the line. "Bytes read past the boundary are put back in front of the stream" is also ambiguous, because they are prepended only in the returned stream, not pushed back into in.

The only caller uses the returned stream, but a future caller who trusts this sentence and keeps reading in (for example the tryWithResource handle) would silently lose up to 8 KiB of events. Something like this would be clearer: "Returns a stream that starts at the first line beginning at or after byte target of in, together with that line's byte offset. in itself is consumed and must not be read afterwards; bytes read past the boundary are prepended to the returned stream."

Verification:

  • Inspection: The scaladoc no longer claims that in itself is positioned at the line start, and its description of the returned stream, the consumption of in and the non-positive-target case matches the implementation.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Updated in cf13cc0 along the lines you suggested. The scaladoc now says that the returned stream starts at the first line beginning at or after target, and that in is consumed and must not be read afterwards, since the bytes read past the boundary are prepended to the returned stream rather than pushed back into in.

// The job description is written as raw UTF-8: skip to the second byte of a character.
val leadByte = bytes.indexWhere(b => (b & 0xF0) == 0xE0)
assert(leadByte > 0)
val conf = createTestConf().set(END_EVENT_REPARSE_CHUNK_SIZE, (bytes.length - leadByte - 1L))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (P3): With bytes.length - leadByte - 1, target is the second byte of U+4E2D (E4 B8 AD), but skipToLineStart now skips only target - 1 bytes, so the raw scan starts on the lead byte E4, a character boundary. The test still catches the original bug (skipping exactly target bytes and then decoding). However, it would also pass with an alignment that decodes UTF-8 from target - 1 before looking for \n, and that alignment fails whenever target - 1 is a continuation byte, which is exactly the case the raw-byte scan is for.

Using bytes.length - leadByte - 2 puts the target on the third byte, so the scan starts on B8. With that change and an updated comment, this test would cover a mid-character start. A couple of multi-byte cases in the skipToLineStart unit test, with the target on each byte of a character, would also pin this down at the helper level.

Verification:

  • Regression: A completed log whose reparse offset falls inside a multi-byte character, with the raw scan starting on a continuation byte, is listed as completed, and a decode-before-scan alignment would make the test fail.
  • Behavior: skipToLineStart returns the first line starting at or after the target, and its UTF-8 byte offset, for targets on the lead, second and third byte of a 3-byte character.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, thanks. Changed in cf13cc0: the target is now the third byte of U+4E2D, so the raw scan starts on the continuation byte B8. I also added cases to the skipToLineStart test with the target on each byte of a 3-byte character and on the following \n.

…multi-byte scan on a continuation byte

Co-authored-by: Claude Code
@viirya

viirya commented Oct 11, 2026

Copy link
Copy Markdown
Member Author

@HyukjinKwon Thanks for raising the listing cost. I measured it and it is real. On the end-event reparse path, with 210 MB of uncompressed event lines (built from spark-events/local-1642039451826 with randomized digits, median of 3 runs, JDK 21):

codec file size reparsed tail Source.getLines() (before this PR) bounded reader (this PR)
lz4 56 MB 155 MB 0.17 s 0.90 s
zstd 16 MB 195 MB 0.17 s 1.05 s

So this PR would make the listing of compressed completed apps about 5-6x slower on this step. The cause is the per-character BufferedReader.read() in boundedLines, which full replays have already been paying since SPARK-59407.

I opened SPARK-60158 / #59360 to read decoded characters in bulk there while keeping the current semantics. With it, a full replay of the same logs takes 0.20-0.22 s instead of 1.10-1.13 s, about the same as Source.getLines().

I'd like to merge #59360 first and this PR after it, so that master never has the slower listing, and to backport the two together. Please hold off on merging this PR until then.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants