Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 12 additions & 5 deletions core/src/main/scala/org/apache/spark/internal/config/History.scala
Original file line number Diff line number Diff line change
Expand Up @@ -165,18 +165,25 @@ private[spark] object History {
.bytesConf(ByteUnit.BYTE)
.createWithDefaultString("1m")

val EVENT_LOG_MAX_LINE_LENGTH_LIMIT: Int = 512 * 1024 * 1024

val EVENT_LOG_MAX_LINE_LENGTH =
ConfigBuilder("spark.history.fs.eventLog.maxLineLength")
.doc("Maximum length of a single event log line during replay. Lines longer than " +
"this are skipped with a warning instead of being read into memory, 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. " +
.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. Buffer " +
"growth, UTF-16 storage and JSON parsing can require several times this limit in heap " +
"space. Buffer capacity grows in steps: reducing 256m to 200m or 150m may not reduce " +
"the retained buffer; use 128m to reach a smaller capacity. Values at or below 0, " +
s"or above $EVENT_LOG_MAX_LINE_LENGTH_LIMIT, use the maximum supported limit of " +
s"$EVENT_LOG_MAX_LINE_LENGTH_LIMIT bytes (512 MiB). This cap avoids JVM array-size " +
"limits but does not guarantee sufficient heap space. " +
"Introduced in 4.3.0; also available in 3.5.10, 4.0.5, 4.1.4 and 4.2.1; and 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.

Following up on my earlier backport comment: v4.2.1-rc1 (cut 2026-09-28 06:11 UTC, about an hour before the latest commit here) already contains SPARK-59407 with the old semantics: .version("4.2.1"), the 512m default, the sb.length() * 2 check, and "Setting this to 0 or a negative value disables the limit". branch-4.2 has already moved to 4.2.2-SNAPSHOT without this fix. My earlier comment only covered branch-4.3, assuming nothing had shipped yet.

If 4.2.1 is finalized from RC1 and this PR is backported afterwards, a patch release (4.2.2) would change the unit, the default and the <= 0 meaning of a released config:

  • Default users: the same for ASCII, but stricter for non-ASCII-heavy lines (e.g. 268M -> 134M chars for 2-byte text, 268M -> ~89.5M chars for 3-byte text).
  • 8k: an ApplicationStart line with 3,000 x U+4E2D (3,085 chars / 9,085 bytes) is replayed by RC1 but skipped here, so the app is missing from the listing.
  • 0 / negative / 3g: RC1 never skips (sb.length() * 2 < Int.MaxValue is always true), while this PR clamps to 512 MiB, so there is no way to replay events above 512 MiB anymore.
  • An explicit 512m (RC1's documented default): RC1 keeps at most 268,435,456 chars and skips U+0101 followed by 600,000,000 ASCII chars with -Xmx1g, while this PR grows the builder to 603,979,774 chars (~1.2 GB as UTF-16) and fails with OutOfMemoryError on UI load or compaction.

Could we get this into branch-4.2 before 4.2.1 is finalized, e.g. by raising it on the RC1 vote so that it lands in RC2? Otherwise, we need a core-migration-guide.md entry for 4.2.1 -> 4.2.2 and should qualify "also available in ... 4.2.1" here and in monitoring.md.

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.

Confirmed that v4.2.1-rc1 contains the old implementation and branch-4.2 is now 4.2.2-SNAPSHOT. I raised this on the Spark dev mailing list so we can decide whether to produce an RC2 with the reviewed fix before finalizing 4.2.1. I am holding off on migration changes pending that discussion.

"all versions after 4.3.0.")
.version("4.3.0")

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.

Note that I found the backport PRs of SPARK-59407 fills this .version field of this config with the version number in each branch, i.e., .version("4.1.4") on branch-4.1 backport, .version("3.5.10") on branch-3.5 backport.

It makes this config look a bit strange as on each version, this config has different version value on different branch. But what .version specifies should be the version that the config is added.

It seems because this is an improvement ticket and adds a new config, but it backports to many previous branches so the config in earlier versions looks strange.

I'm not sure how to deal with it when to backport this fix to these branches. Just backport this so all branches have "4.3.0" in .version of this config? Or leave the strange version unchanged? 🤔

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.

@dongjoon-hyun @HyukjinKwon What do you think?

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.

I think we have done this with 3.5.10 in this case so far IIRC

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.

Do you mean to set .version("3.5.10") for all branches?

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.

Actually, I'm not sure about that. Both different versions and 3.5.10 sounds misleading to me. However, if there is any precedence, I don't care.

Note that I found the backport PRs of SPARK-59407 fills this .version field of this config with the version number in each branch, i.e., .version("4.1.4") on branch-4.1 backport, .version("3.5.10") on branch-3.5 backport.

cc @holdenk for looping the author of that backport.

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.

If not 3.5.10, what version string we should use in these branches? 4.3.0 as the doc said that it is introduced in 4.3.0?

.withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE)
.bytesConf(ByteUnit.BYTE)
.createWithDefaultString("512m")
.createWithDefaultString("256m")

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.

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.3 does, so this needs to land on branch-4.3 before the next 4.3.0 RC. Otherwise 4.3.0 ships the old semantics (e.g. 8k means about 4096 chars and <= 0 means 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 new case jpe: JsonProcessingException block with log"..."/MDC merges 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.

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.

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.


private[spark] val EVENT_LOG_ROLLING_MAX_FILES_TO_RETAIN =
ConfigBuilder("spark.history.fs.eventLog.rolling.maxFilesToRetain")
Expand Down
143 changes: 101 additions & 42 deletions core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import java.nio.charset.{CodingErrorAction, StandardCharsets}

import scala.annotation.tailrec

import com.fasterxml.jackson.core.JsonParseException
import com.fasterxml.jackson.core.{JsonParseException, JsonProcessingException}
import com.fasterxml.jackson.databind.exc.UnrecognizedPropertyException

import org.apache.spark.SparkConf
Expand All @@ -35,15 +35,21 @@ import org.apache.spark.util.JsonProtocol
/**
* A SparkListenerBus that can be used to replay events from serialized event data.
*
* @param maxLineLength Maximum number of byes (1/2 chars) of a single event log line that will be
* materialized during replay. Longer lines are drained, skipped and
* logged, bounding the memory replay can use when an event log is
* corrupt or unexpectedly large.
* @param maxLineLength Maximum UTF-8 byte length of a single event log line, excluding its line

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: Could the @param mention the normalization, e.g. "Values <= 0 or above ReplayListenerBus.MAX_LINE_LENGTH (512 MiB) use MAX_LINE_LENGTH"? Otherwise a caller passing 0 (RC1's "disables the limit" semantics) or 1 << 30 silently gets a 512 MiB limit.

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.

The constructor scaladoc now describes the normalization of non-positive and above-maximum values to the 512 MiB limit.

* ending. Longer lines are drained, skipped and logged, bounding the
* memory replay can use when an event log is corrupt or unexpectedly large.
* Values at or below zero or above MAX_LINE_LENGTH use MAX_LINE_LENGTH
* (512 MiB).
*/
private[spark] class ReplayListenerBus(

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.

Binary compatibility (pre-existing from SPARK-59407, but this PR is a good place to fix it before the maintenance releases): every released version (3.5.9, 4.0.4, 4.1.3, 4.2.0) declares class ReplayListenerBus extends SparkListenerBus with Logging, which has a public no-arg constructor. With (maxLineLength: Int = DEFAULT_MAX_LINE_LENGTH), scalac only emits <init>(I)V; default arguments don't generate a no-arg constructor.

The class is private[spark] and excluded from MiMa, but event-log tools already use it and would break on 3.5.10 / 4.0.5 / 4.1.4 / 4.2.1:

The reflective cases fail even after recompiling. Could we add an auxiliary constructor?

  def this() = this(ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH)

It compiles on Scala 2.12 and 2.13 with or without the default argument, new ReplayListenerBus() is not ambiguous, and all in-repo call shapes still compile, including new ReplayListenerBus(maxLineLength = ...) and the anonymous subclass in StreamingQueryListenerSuite. If 4.2.1 ships from RC1, keeping the default argument as well also preserves $lessinit$greater$default$1 for code compiled against 4.2.1.

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.

Confirmed with the compiled class: the no-argument constructor was missing. I restored it while retaining the default argument, and added a reflection regression test. The compiled class now exposes both constructors and the Scala default-argument method.

maxLineLength: Int = ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH)
extends SparkListenerBus with Logging {

def this() = this(ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH)

private[scheduler] val effectiveMaxLineLength =
ReplayListenerBus.normalizeMaxLineLength(maxLineLength)

/**
* Replay each event in the order maintained in the given stream. The stream is expected to
* contain one JSON-encoded SparkListenerEvent per line.
Expand All @@ -67,22 +73,26 @@ private[spark] class ReplayListenerBus(
maybeTruncated: Boolean = false,
eventsFilter: ReplayEventsFilter = SELECT_ALL_FILTER): Boolean = {
val lines = boundedLines(logData, sourceName)
replay(lines, sourceName, maybeTruncated, eventsFilter)
replayEntries(lines, sourceName, maybeTruncated, eventsFilter)
}

/**
* Reads '\n'-terminated lines like Source.getLines(), but never materializes more than
* `maxLineLength` bytes of a single line. An over-long line is drained and skipped
* with a warning instead of being turned into a String.
* Reads '\n'-terminated lines and retains their original zero-based indices. The limit
* measures UTF-8 content bytes, excluding the line ending, rather than JVM heap usage.
* The character buffer is bounded by the limit plus one possible trailing CR. An over-long
* line is drained and skipped with a warning instead of being turned into a String.
*/
private def boundedLines(logData: InputStream, sourceName: String): Iterator[String] = {
private def boundedLines(

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.

Since boundedLines is private, compaction reads the same files with two different readers: EventLogFileCompactor.initializeBuilders replays through this bounded reader, while EventFilter.applyFilterToFile (EventFilter.scala L79) rewrites them with an unbounded Source.getLines(). If a skipped over-limit line is the start event of a still-live entity (e.g. StageSubmitted of a running stage, JobStart of a running job, SQLExecutionStart of a running execution without jobs), the builders never mark it live, the rewrite materializes the line and rejects it together with its dependent task events, and cleanupCompactedFiles deletes the originals, so raising the limit later cannot recover them. A port of the builder/filter logic lost 23 events with a 1 MiB limit and 2 MiB lines. The inconsistency came with SPARK-59407 and this PR only moves the threshold (e.g. CJK-heavy lines between limit/3 and limit/2 chars are newly skipped), but sharing one bounded reader for both passes, and passing over-long lines through unchanged in the rewrite, would make them consistent.

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.

We will address compaction reader and retention consistency in a separate follow-up, including reproducing the possible loss of live-entity events. That needs its own retention tests and a decision about handling oversized start events across both passes. This PR does not change the compaction rewrite path.

logData: InputStream,
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)

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.

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.

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.

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.

.onUnmappableCharacter(CodingErrorAction.REPORT)
val reader = new BufferedReader(new InputStreamReader(logData, decoder))

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.

A thought for the block-read follow-up, not a blocker here: the limit is a UTF-8 byte count, but the reader decodes to UTF-16 first and re-derives the byte width per char. Since 0x0A never appears inside a multi-byte UTF-8 sequence, splitting lines on raw bytes into a buffer capped at limit + 1 bytes, and strictly decoding only the lines that are kept, makes the unit exact by construction. Skipped bytes can still be streamed through the strict decoder to keep the "malformed UTF-8 in skipped lines" contract. It would also:

  • bound a skipped non-Latin-1 line to limit + 1 bytes instead of a UTF-16 builder (~604 MB plus growth at 256m),
  • let a MalformedInputException be attributed to a physical line (today the decoder's 8 KiB read-ahead can throw before the preceding lines are delivered),
  • provide the raw-byte handling that the tail-reparse (partial first line) and truncated-UTF-8-at-EOF follow-ups need.

A quick Java port matched this reader on all the grids in this PR plus random and malformed inputs. The reader alone ran at ~2,450 / ~310 MB/s (ASCII / CJK), vs ~77 / ~91 MB/s here and ~1,140 / ~215 MB/s for a bulk char[] reader. It still needs the CR reservation, a cap (around the 1 GiB UTF-16 String limit), and a heap-oriented default.

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.

We will evaluate raw-byte splitting in the reader follow-up, preserving strict validation of skipped content and covering tail alignment and incomplete UTF-8 at EOF. I have not independently reproduced the benchmark numbers.

new Iterator[String] {
private var nextLine: String = _
new Iterator[(String, Int)] {
private var nextLine: (String, Int) = _
private var lineIndex = 0
private var lineFetched = false
private var warned = false

Expand All @@ -94,7 +104,7 @@ private[spark] class ReplayListenerBus(
nextLine != null
}

override def next(): String = {
override def next(): (String, Int) = {
if (!hasNext) {
throw new NoSuchElementException("No more lines")
}
Expand All @@ -104,35 +114,30 @@ private[spark] class ReplayListenerBus(
line
}

@tailrec private def fetchLine(): String = {
val sb = new java.lang.StringBuilder()
var overLong = false
@tailrec private def fetchLine(): (String, Int) = {
val buffer = new BoundedLineBuffer(effectiveMaxLineLength)

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.

For the block-read follow-up: a new BoundedLineBuffer (and a 16-char StringBuilder regrown by doubling) per line, including skipped lines, allocates ~4.3 bytes per input byte on the event logs under core/src/test/resources/spark-events, vs ~1.7 with one buffer per iterator that is reset for each line (Jackson readTree alone is ~3.9). Today the per-char BufferedReader.read() lock hides the CPU cost, but with a bulk reader it becomes ~10% of the read loop. Most of this is inherited from SPARK-59407 (this PR adds ~40 B per line), so it's fine to handle it there. A reused buffer would need a capacity cap so that it doesn't pin a huge array for the rest of the file.

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 to evaluate buffer reuse in the reader follow-up, including a retention cap after unusually large lines. I have not independently reproduced the allocation or throughput measurements.

var c = reader.read()
if (c == -1) {
null
} else {
val index = lineIndex

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: This hand-rolls what zipWithIndex already does at L151. If the reader emits one Option per physical line, both overloads number lines the same way, and the manual counter, the tuple-typed state and the @tailrec recursion go away:

    var warned = false
    Iterator.continually(reader.read()).takeWhile(_ != -1).map { first =>
      val buffer = new BoundedLineBuffer(effectiveMaxLineLength)
      var c = first
      while (c != -1 && c != '\n') {
        buffer.append(c.toChar)
        c = reader.read()
      }
      val maybeLine = buffer.result()
      if (maybeLine.isEmpty && !warned) {
        logWarning(...) // unchanged
        warned = true
      }
      maybeLine
    }.zipWithIndex.collect { case (Some(line), index) => (line, index) }

A NextIterator[Option[String]] also works if relying on takeWhile's one-element look-ahead feels too implicit. I checked that this behaves the same as the current iterator on Scala 2.12 and 2.13, including the look-ahead from lineEntries.hasNext in the JsonParseException handler and exceptions thrown mid-line. Feel free to skip it if you'd rather keep the diff against SPARK-59407 small.

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.

I would keep the current iterator in this fix to limit the backport diff. We can revisit the iterator structure with the reader follow-up.

lineIndex += 1
while (c != -1 && c != '\n') {
if (sb.length() * 2 < maxLineLength) {
sb.append(c.toChar)
} else {
overLong = true
}
buffer.append(c.toChar)
c = reader.read()

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.

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?

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 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.

}
if (overLong) {
val maybeLine = buffer.result()
if (maybeLine.isEmpty) {
if (!warned) {
logWarning(log"Skipped event log lines longer than " +

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.

Since the physical index is available here now, could the warning include it (e.g. index + 1 of the first skipped line)? Today a replay can return true after dropping several lines, and the only trace is one "Skipped event log lines longer than 268435456 bytes in " per file, without a location or a count; later skips are not logged at any level. replayEntries already uses one WARN plus a per-occurrence logDebug for unrecognized events (L176-L183), which would fit here too. A count at EOF alone wouldn't be enough, because listing replays often stop early with HaltReplayException.

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 the first skipped physical line number to the warning and a DEBUG message for every skipped line. A regression test checks one warning and the locations of both skipped lines.

log"${MDC(MAX_SIZE, maxLineLength)} characters in " +
log"${MDC(FILE_NAME, sourceName)}")
log"${MDC(MAX_SIZE, effectiveMaxLineLength)} bytes in " +
log"${MDC(FILE_NAME, sourceName)}; first skipped line: ${MDC(LINE_NUM, index + 1)}")
warned = true
}
logDebug(s"Skipped event log line ${index + 1} in $sourceName")
fetchLine()
} else {
// Handle CRLF line endings like Source.getLines() does.
if (sb.length() > 0 && sb.charAt(sb.length() - 1) == '\r') {
sb.setLength(sb.length() - 1)
}
sb.toString
(maybeLine.get, index)
}
}
}
Expand All @@ -148,15 +153,21 @@ private[spark] class ReplayListenerBus(
sourceName: String,
maybeTruncated: Boolean,
eventsFilter: ReplayEventsFilter): Boolean = {
replayEntries(lines.zipWithIndex, sourceName, maybeTruncated, 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.

This overload is also what the end-event reparse in FsHistoryProvider uses (FsHistoryProvider.scala L1109-L1127), fed by an unbounded Source.fromInputStream(in)(Codec.UTF8).getLines(). So on that path there is no maxLineLength and line numbers are relative to the skip point. With the defaults (compressed event logs, endEventReparseChunkSize=1m), completed apps always take this path, and because target = lastFile.getLen - reparseChunkSize is based on the compressed length while skip runs on the decompressed stream, most of the log is read there. In a local repro (156.6 MB decompressed / 16.4 MB compressed), 141 MB went through getLines(), and a 12.6 MB line over a 4 MiB limit was fully materialized. This predates the PR, but since the PR defines the limit and line-number semantics, it may be worth reusing the bounded reader there (skip, drop the partial first line, then replay through boundedLines) or filing a follow-up.

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.

We will handle this in a separate follow-up. The tail-reparse reader and its offset/partial-line semantics are independent of the two issues addressed here. This PR leaves that path unchanged; the follow-up will address the bypass of the bounded reader and clarify line-number semantics after skipping into a file.

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.

Two notes for the tail-reparse follow-up (FsHistoryProvider L1104-L1127), which calls this overload:

  1. There, lines starts 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 truncated SparkListenerApplicationEnd line at physical line 12,004 of a completed app is reported as at 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 existing Malformed line #N and 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)".
  2. Pre-existing: if target lands 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) throws MalformedInputException before replay is called. mergeApplicationListing only logs it, and since the placeholder LogInfo was already written with the same file size, shouldReloadLog never retries it, so the completed app stays missing from the listing (and target is deterministic, so a restart fails the same way). Replaying through boundedLines after 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 \n before decoding. I reproduced it at the unit level, not yet on a real SHS.

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.

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.

}

private def replayEntries(

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.

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.

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.

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.

lines: Iterator[(String, Int)],
sourceName: String,
maybeTruncated: Boolean,
eventsFilter: ReplayEventsFilter): Boolean = {
var currentLine: String = null
var lineNumber: Int = 0
val unrecognizedEvents = new scala.collection.mutable.HashSet[String]
val unrecognizedProperties = new scala.collection.mutable.HashSet[String]

try {
val lineEntries = lines
.zipWithIndex
.filter { case (line, _) => eventsFilter(line) }
val lineEntries = lines.filter { case (line, _) => eventsFilter(line) }

while (lineEntries.hasNext) {
try {
Expand Down Expand Up @@ -202,11 +213,18 @@ private[spark] class ReplayListenerBus(
// Just stop replay.
false
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)}")
throw jpe
case ioe: IOException =>
throw ioe
case e: Exception =>

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.

Pre-existing and optional: the new JSON branch above logs only the source and line, but this generic branch still logs the entire line (Malformed line #...: ${MDC(LINE, currentLine)} at L219), which can be up to 256 MiB by default (512 MiB max). A large line that is valid JSON but semantically wrong ends up here, e.g. a JobStart whose Properties contain a numeric value fails extractString's require with IllegalArgumentException. With the same code in Spark 3.5.0, a 33,554,544-char line produced a single 33,554,606-char ERROR record, and building the message temporarily needs about twice the line size (the log"..." builder doubles when the trailing \n is appended). Since this PR is about bounding replay memory, could we log a bounded prefix plus the line length instead? The Malformed line #4: {} assertion would still pass.

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.

The malformed-line diagnostic now includes at most the first 1,024 characters and reports the full character count. A regression test uses valid JSON with a numeric property value to exercise the generic exception path.

logError(log"Exception parsing Spark event log: ${MDC(PATH, sourceName)}", e)
logError(log"Malformed line #${MDC(LINE_NUM, lineNumber)}: ${MDC(LINE, currentLine)}\n")
val prefix = Option(currentLine).map(_.take(1024)).orNull
val length = Option(currentLine).fold(0)(_.length)
logError(log"Malformed line #${MDC(LINE_NUM, lineNumber)}: ${MDC(LINE, prefix)} " +
log"(line length: ${MDC(SIZE, length)} characters; showing at most 1024)\n")
false
}
}
Expand All @@ -225,21 +243,62 @@ private[spark] class HaltReplayException extends RuntimeException

private[spark] object ReplayListenerBus {

/**
* Default per-line cap during replay: far above any legitimate event line, so replay
* memory stays bounded even for corrupt logs. Matches the default of
* spark.history.fs.eventLog.maxLineLength.
*/
val DEFAULT_MAX_LINE_LENGTH: Int = 512 * 1024 * 1024
/** Default UTF-8 content limit shared with the history server configuration. */
val DEFAULT_MAX_LINE_LENGTH: Int = History.EVENT_LOG_MAX_LINE_LENGTH.defaultValue.get.toInt

// Bound StringBuilder growth so UTF-16 inflation stays below the JVM array-size limit.
val MAX_LINE_LENGTH: Int = History.EVENT_LOG_MAX_LINE_LENGTH_LIMIT

/** Resolves the replay line-length cap from configuration; <= 0 disables the cap. */
/** Resolves the byte limit, using 512 MiB for non-positive or larger values. */
def maxLineLength(conf: SparkConf): Int = {
val configured = conf.get(History.EVENT_LOG_MAX_LINE_LENGTH)
if (configured <= 0 || configured > Int.MaxValue) Int.MaxValue else configured.toInt
normalizeMaxLineLength(conf.get(History.EVENT_LOG_MAX_LINE_LENGTH))
}

private def normalizeMaxLineLength(configured: Long): Int = {
if (configured <= 0 || configured > MAX_LINE_LENGTH) MAX_LINE_LENGTH else configured.toInt
}

type ReplayEventsFilter = (String) => Boolean

// utility filter that selects all event logs during replay
val SELECT_ALL_FILTER: ReplayEventsFilter = { (eventString: String) => true }

/** Accumulates a bounded prefix of one decoded line, allowing a possible trailing CR. */
private[scheduler] class BoundedLineBuffer(maxLineLength: Int) {
private val buffer = new java.lang.StringBuilder()

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.

FYI, because this StringBuilder starts at 16 and grows as 2c + 2, the retained buffer is a step function of the limit rather than proportional to it. For ASCII-heavy lines, every limit from 150,994,942 to 301,989,885 bytes (~144m to ~288m) ends with the same 301,989,886-char buffer as the 256m default, and ~288m to 512m all end with 603,979,774 chars.

So if an operator hits an OOM while a long line is drained and lowers the limit to 200m or 150m, nothing changes. With -Xmx800m, U+0101 followed by 300,000,000 ASCII chars failed at the same point (length 150,994,942, when the 302 MB UTF-16 array is copied into a 604 MB one) for 256m, 200m and 150m, while 143m and 128m skipped it. Could the docs mention that only limits below ~144m (e.g. 128m) reduce the buffer, or could the growth be capped at bufferLimit?

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.

Confirmed the capacity steps from the JDK growth rule. I documented that lowering 256m to 200m or 150m may leave the same buffer capacity, with 128m as an example of a smaller step. Changing buffer growth remains part of the reader follow-up.

private var byteLength = 0L
private val bufferLimit = maxLineLength.toLong + 1

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.

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 with UTF16 String size is 1207959550, should be less than 1073741823. This also hits lines within the limit: with 600m, 620,000,000 x followed 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 billion x at 1g throws Requested array size exceeds VM limit, while the current code skips it.
  • <= 0 and values above Int.MaxValue fail 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).

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.

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.


def length: Int = buffer.length()

def capacity: Int = buffer.capacity()

def append(c: Char): Unit = {
if (byteLength <= bufferLimit) {
// The strict decoder emits valid surrogate pairs: count two bytes per half.
byteLength += (if (c < 0x80) 1 else if (c < 0x800 || Character.isSurrogate(c)) 2 else 3)
// Reserve one extra byte for a possible CR in the line ending.
if (byteLength <= bufferLimit) {

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.

Once byteLength exceeds bufferLimit, the line is certain to be skipped, but the retained prefix stays reachable until the \n is read. With the default 256m, that is ~302 MB (Latin-1) or ~604 MB (UTF-16, i.e. any char >= U+0100 early in the line) of dead heap for the whole drain (~13 s per GB at the current read rate), overlapping with concurrent UI loads and compaction. In a -Xmx1g model, bus A draining U+0101 + ~868 MB of ASCII made bus B fail with OutOfMemoryError on a 200 MB line, and releasing the buffer on overflow made it pass.

This is pre-existing from SPARK-59407, but it's a small fix:

        if (byteLength <= bufferLimit) {
          buffer.append(c)
        } else {
          // The line will be skipped; release the retained prefix while draining.
          buffer.setLength(0)
          buffer.trimToSize()
        }

After overflow byteLength >= maxLineLength + 2, so result() still returns None regardless of a trailing CR. It doesn't lower the peak, only how long the memory is held. The assertion at ReplayListenerSuite L164 (buffer.length == limit + 1) would become == 0.

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.

The retained prefix is now cleared and trimmed once overflow makes the line unconditionally skippable. The buffer test checks both zero length and zero capacity after overflow; this reduces retention during draining, not the allocation peak.

buffer.append(c)
} else {
// The line will be skipped; release the retained prefix while draining.
buffer.setLength(0)
buffer.trimToSize()
}
}
}

def result(): Option[String] = {
val trailingCR = length > 0 && buffer.charAt(length - 1) == '\r'
val contentLength = byteLength - (if (trailingCR) 1 else 0)
if (contentLength > maxLineLength) {
None
} else if (trailingCR) {
Some(buffer.substring(0, length - 1))
} else {
Some(buffer.toString)
}
}
}

}
Loading