-
Notifications
You must be signed in to change notification settings - Fork 29.4k
[SPARK-59804][UI] Fix UTF-8 byte limits and physical line numbers during event log replay #59074
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
1df219a
ab7c85d
aaee04e
691a1d7
1167bfa
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 " + | ||
| "all versions after 4.3.0.") | ||
| .version("4.3.0") | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Note that I found the backport PRs of SPARK-59407 fills this It makes this config look a bit strange as on each version, this config has different version value on different branch. But what 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. @dongjoon-hyun @HyukjinKwon What do you think?
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Do you mean to set
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Actually, I'm not sure about that. Both different versions and
cc @holdenk for looping the author of that backport.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. If not |
||
| .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) | ||
| .bytesConf(ByteUnit.BYTE) | ||
| .createWithDefaultString("512m") | ||
| .createWithDefaultString("256m") | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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") | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 | ||
|
|
@@ -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 | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: Could the
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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( | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 The class is
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,
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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. | ||
|
|
@@ -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( | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Since
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Pre-existing and preserved by this PR ( 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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)) | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
|
||
|
|
@@ -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") | ||
| } | ||
|
|
@@ -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) | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. For the block-read follow-up: a new
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: This hand-rolls what 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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() | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Performance (pre-existing since SPARK-59407, kept by this rewrite):
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 " + | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Since the physical
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
| } | ||
| } | ||
| } | ||
|
|
@@ -148,15 +153,21 @@ private[spark] class ReplayListenerBus( | |
| sourceName: String, | ||
| maybeTruncated: Boolean, | ||
| eventsFilter: ReplayEventsFilter): Boolean = { | ||
| replayEntries(lines.zipWithIndex, sourceName, maybeTruncated, eventsFilter) | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This overload is also what the end-event reparse in
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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( | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The physical line number only surfaces for non-IOException failures.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Confirmed and fixed in ab7c85d. 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 { | ||
|
|
@@ -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 => | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 (
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
| } | ||
| } | ||
|
|
@@ -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() | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. FYI, because this So if an operator hits an OOM while a long line is drained and lowers the limit to
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This retains up to
On the listing path, the Could we clamp the effective limit in
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Once 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
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
| } | ||
| } | ||
| } | ||
|
|
||
| } | ||
There was a problem hiding this comment.
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"), the512mdefault, thesb.length() * 2check, and "Setting this to 0 or a negative value disables the limit".branch-4.2has already moved to4.2.2-SNAPSHOTwithout this fix. My earlier comment only coveredbranch-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
<= 0meaning of a released config:8k: anApplicationStartline 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.MaxValueis always true), while this PR clamps to 512 MiB, so there is no way to replay events above 512 MiB anymore.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 withOutOfMemoryErroron UI load or compaction.Could we get this into
branch-4.2before 4.2.1 is finalized, e.g. by raising it on the RC1 vote so that it lands in RC2? Otherwise, we need acore-migration-guide.mdentry for 4.2.1 -> 4.2.2 and should qualify "also available in ... 4.2.1" here and inmonitoring.md.There was a problem hiding this comment.
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.