From 1df219aec39bd795de704908b6f089fc50bc856f Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sat, 26 Sep 2026 17:12:43 -0700 Subject: [PATCH 1/5] [SPARK-59804][UI] Fix UTF-8 byte limits and physical line numbers during event log replay --- .../spark/internal/config/History.scala | 4 +- .../spark/scheduler/ReplayListenerBus.scala | 82 +++++++++++++------ .../spark/scheduler/ReplayListenerSuite.scala | 79 ++++++++++++++++++ docs/monitoring.md | 4 +- 4 files changed, 138 insertions(+), 31 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/internal/config/History.scala b/core/src/main/scala/org/apache/spark/internal/config/History.scala index 7353481b6582f..d56ba0cbe59ac 100644 --- a/core/src/main/scala/org/apache/spark/internal/config/History.scala +++ b/core/src/main/scala/org/apache/spark/internal/config/History.scala @@ -167,8 +167,8 @@ private[spark] object History { 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 " + + .doc("Maximum UTF-8 byte length of a single event log line during replay, excluding " + + "the line ending. Longer lines are skipped with a warning, bounding the " + "memory replay can use when an event log is corrupt or unexpectedly large. Setting " + "this to 0 or a negative value disables the limit. " + "Introduced in 4.3.0; also available in 3.5.10, 4.0.5, 4.1.4 and 4.2.1; and in " + diff --git a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala index 1cd00facb7396..6a286899c8fa1 100644 --- a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala +++ b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala @@ -35,10 +35,9 @@ 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 + * ending. Longer lines are drained, skipped and logged, bounding the + * memory replay can use when an event log is corrupt or unexpectedly large. */ private[spark] class ReplayListenerBus( maxLineLength: Int = ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) @@ -67,22 +66,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( + 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) .onUnmappableCharacter(CodingErrorAction.REPORT) val reader = new BufferedReader(new InputStreamReader(logData, decoder)) - 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 +97,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 +107,54 @@ private[spark] class ReplayListenerBus( line } - @tailrec private def fetchLine(): String = { + @tailrec private def fetchLine(): (String, Int) = { val sb = new java.lang.StringBuilder() var overLong = false + var byteLength = 0L var c = reader.read() if (c == -1) { null } else { + val index = lineIndex + lineIndex += 1 while (c != -1 && c != '\n') { - if (sb.length() * 2 < maxLineLength) { - sb.append(c.toChar) - } else { - overLong = true + if (!overLong) { + // The strict decoder guarantees valid surrogate pairs. Count their four UTF-8 + // bytes on the high surrogate and zero on the low surrogate. + byteLength += (if (c < 0x80) { + 1 + } else if (c < 0x800) { + 2 + } else if (Character.isHighSurrogate(c.toChar)) { + 4 + } else if (Character.isLowSurrogate(c.toChar)) { + 0 + } else { + 3 + }) + // Allow one extra byte until we know whether the line ends in CRLF. + if (byteLength <= maxLineLength.toLong + 1) { + sb.append(c.toChar) + } else { + overLong = true + } } c = reader.read() } - if (overLong) { + if (sb.length() > 0 && sb.charAt(sb.length() - 1) == '\r') { + sb.setLength(sb.length() - 1) + byteLength -= 1 + } + if (overLong || byteLength > maxLineLength) { if (!warned) { logWarning(log"Skipped event log lines longer than " + - log"${MDC(MAX_SIZE, maxLineLength)} characters in " + + log"${MDC(MAX_SIZE, maxLineLength)} bytes in " + log"${MDC(FILE_NAME, sourceName)}") warned = true } 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 + (sb.toString, index) } } } @@ -148,15 +170,21 @@ private[spark] class ReplayListenerBus( sourceName: String, maybeTruncated: Boolean, eventsFilter: ReplayEventsFilter): Boolean = { + replayEntries(lines.zipWithIndex, sourceName, maybeTruncated, eventsFilter) + } + + private def replayEntries( + 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 { diff --git a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala index f4652c23f63e2..b8855a9633ea7 100644 --- a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala @@ -109,6 +109,85 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp assert(eventMonster.loggedEvents(1) === JsonProtocol.sparkEventToJsonString(applicationEnd)) } + test("Replay line limit uses UTF-8 bytes") { + val start = SparkListenerApplicationStart("x" * 6000, None, 125L, "user", None) + val json = JsonProtocol.sparkEventToJsonString(start) + assert(json.getBytes(StandardCharsets.UTF_8).length < 8 * 1024) + val listener = new EventBufferingListener + val bus = new ReplayListenerBus(8 * 1024) + bus.addListener(listener) + val input = new ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8)) + assert(bus.replay(input, "ascii")) + assert(listener.loggedEvents.toSeq == Seq(json)) + } + + test("Replay line limit handles UTF-8 boundaries and line endings") { + val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + // scalastyle:off nonascii + val names = Seq("x", "\u00e9", "\u4e2d", "\ud83d\ude00") + // scalastyle:on nonascii + for { + name <- names + repetitions <- Seq(100, 4097) + ending <- Seq("\n", "\r\n", "") + excess <- Seq(0, 1) + } { + val start = SparkListenerApplicationStart(name * repetitions, None, 125L, "user", None) + val json = JsonProtocol.sparkEventToJsonString(start) + val limit = json.getBytes(StandardCharsets.UTF_8).length - excess + val listener = new EventBufferingListener + val bus = new ReplayListenerBus(limit) + bus.addListener(listener) + // Put a short event first, including when the final line has no terminator. + val input = new ByteArrayInputStream((end + "\n" + json + ending) + .getBytes(StandardCharsets.UTF_8)) + assert(bus.replay(input, "utf8")) + val expected = if (excess == 0) Seq(end, json) else Seq(end) + assert(listener.loggedEvents.toSeq == expected, s"name=$name ending=$ending excess=$excess") + } + } + + test("Replay preserves physical line numbers after skipping long lines") { + val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + val appender = new LogAppender + val bus = new ReplayListenerBus(1024) + val input = new ByteArrayInputStream( + (end + "\n" + "x" * 2048 + "\n" + "x" * 2048 + "\n{}\n") + .getBytes(StandardCharsets.UTF_8)) + withLogAppender(appender) { + assert(!bus.replay(input, "line-numbers")) + } + assert(appender.loggingEvents.exists(_.getMessage.getFormattedMessage + .contains("Malformed line #4: {}"))) + } + + test("Replay preserves line numbers for truncated logs and filtered iterators") { + val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + val truncated = "{\"Event\":" + val appender = new LogAppender + val bus = new ReplayListenerBus(1024) + val input = new ByteArrayInputStream( + ("x" * 2048 + "\n" + end + "\n" + truncated).getBytes(StandardCharsets.UTF_8)) + withLogAppender(appender) { + assert(bus.replay(input, "truncated", maybeTruncated = true, + eventsFilter = _ != end)) + assert(bus.replay(Iterator(end, end, truncated), "iterator", maybeTruncated = true, + eventsFilter = _ != end)) + } + val warnings = appender.loggingEvents.map(_.getMessage.getFormattedMessage) + .filter(_.contains("Got JsonParseException")) + assert(warnings.size == 2) + assert(warnings.forall(_.contains("at line 3,"))) + } + + test("Replay still rejects malformed UTF-8 in skipped lines") { + val bytes = Array.fill[Byte](20 * 1024)('x'.toByte) ++ Array(0xff.toByte, '\n'.toByte) + val bus = new ReplayListenerBus(1024) + intercept[java.nio.charset.MalformedInputException] { + bus.replay(new ByteArrayInputStream(bytes), "invalid-utf8") + } + } + /** * Test replaying compressed spark history file that internally throws an EOFException. To * avoid sensitivity to the compression specifics the test forces an EOFException to occur diff --git a/docs/monitoring.md b/docs/monitoring.md index 83fd85065f446..c81b6569cfa09 100644 --- a/docs/monitoring.md +++ b/docs/monitoring.md @@ -442,8 +442,8 @@ Security options for the Spark History Server are covered more detail in the spark.history.fs.eventLog.maxLineLength 512m - 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, which bounds the memory replay + Maximum UTF-8 byte length of a single event log line during replay, excluding the line + ending. Longer lines are skipped with a warning, which bounds 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.
Introduced in 4.3.0; also available in 3.5.10, 4.0.5, 4.1.4 and 4.2.1; and in all From ab7c85da8f6648b36853ff84e13e29ea0388221b Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 27 Sep 2026 16:25:30 -0700 Subject: [PATCH 2/5] [SPARK-59804][UI] Address replay limit and diagnostics review --- .../spark/internal/config/History.scala | 8 +- .../spark/scheduler/ReplayListenerBus.scala | 83 ++++++++------ .../spark/scheduler/ReplayListenerSuite.scala | 103 +++++++++++++++--- docs/monitoring.md | 7 +- 4 files changed, 141 insertions(+), 60 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/internal/config/History.scala b/core/src/main/scala/org/apache/spark/internal/config/History.scala index d56ba0cbe59ac..83a28429fa79b 100644 --- a/core/src/main/scala/org/apache/spark/internal/config/History.scala +++ b/core/src/main/scala/org/apache/spark/internal/config/History.scala @@ -169,14 +169,16 @@ private[spark] object History { ConfigBuilder("spark.history.fs.eventLog.maxLineLength") .doc("Maximum UTF-8 byte length of a single event log line during replay, excluding " + "the line ending. Longer lines are skipped with a warning, bounding the " + - "memory replay can use when an event log is corrupt or unexpectedly large. Setting " + - "this to 0 or a negative value disables the limit. " + + "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. Values at or below 0, or above 2147483647, use the maximum supported limit of " + + "2147483647 bytes. " + "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") .withBindingPolicy(ConfigBindingPolicy.NOT_APPLICABLE) .bytesConf(ByteUnit.BYTE) - .createWithDefaultString("512m") + .createWithDefaultString("256m") private[spark] val EVENT_LOG_ROLLING_MAX_FILES_TO_RETAIN = ConfigBuilder("spark.history.fs.eventLog.rolling.maxFilesToRetain") diff --git a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala index 6a286899c8fa1..db8984f41e546 100644 --- a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala +++ b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala @@ -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 @@ -108,9 +108,7 @@ private[spark] class ReplayListenerBus( } @tailrec private def fetchLine(): (String, Int) = { - val sb = new java.lang.StringBuilder() - var overLong = false - var byteLength = 0L + val buffer = new BoundedLineBuffer(maxLineLength) var c = reader.read() if (c == -1) { null @@ -118,34 +116,11 @@ private[spark] class ReplayListenerBus( val index = lineIndex lineIndex += 1 while (c != -1 && c != '\n') { - if (!overLong) { - // The strict decoder guarantees valid surrogate pairs. Count their four UTF-8 - // bytes on the high surrogate and zero on the low surrogate. - byteLength += (if (c < 0x80) { - 1 - } else if (c < 0x800) { - 2 - } else if (Character.isHighSurrogate(c.toChar)) { - 4 - } else if (Character.isLowSurrogate(c.toChar)) { - 0 - } else { - 3 - }) - // Allow one extra byte until we know whether the line ends in CRLF. - if (byteLength <= maxLineLength.toLong + 1) { - sb.append(c.toChar) - } else { - overLong = true - } - } + buffer.append(c.toChar) c = reader.read() } - if (sb.length() > 0 && sb.charAt(sb.length() - 1) == '\r') { - sb.setLength(sb.length() - 1) - byteLength -= 1 - } - if (overLong || byteLength > maxLineLength) { + val line = buffer.result() + if (line.isEmpty) { if (!warned) { logWarning(log"Skipped event log lines longer than " + log"${MDC(MAX_SIZE, maxLineLength)} bytes in " + @@ -154,7 +129,7 @@ private[spark] class ReplayListenerBus( } fetchLine() } else { - (sb.toString, index) + (line.get, index) } } } @@ -230,6 +205,10 @@ 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)}", jpe) + throw jpe case ioe: IOException => throw ioe case e: Exception => @@ -254,13 +233,12 @@ 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. + * Default UTF-8 line-length cap, matching spark.history.fs.eventLog.maxLineLength. + * Keep the maximum buffered character count close to the original 512 MiB / 2 cap. */ - val DEFAULT_MAX_LINE_LENGTH: Int = 512 * 1024 * 1024 + val DEFAULT_MAX_LINE_LENGTH: Int = 256 * 1024 * 1024 - /** Resolves the replay line-length cap from configuration; <= 0 disables the cap. */ + /** Resolves the byte limit, using Int.MaxValue 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 @@ -270,4 +248,37 @@ private[spark] object ReplayListenerBus { // 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() + private var byteLength = 0L + private val bufferLimit = maxLineLength.toLong + 1 + + def length: Int = buffer.length() + + 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) { + buffer.append(c) + } + } + } + + 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) + } + } + } + } diff --git a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala index b8855a9633ea7..a11db118870c9 100644 --- a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala @@ -23,6 +23,8 @@ import java.util.concurrent.atomic.AtomicInteger import scala.collection.mutable.ArrayBuffer +import com.fasterxml.jackson.core.{JsonParseException, JsonProcessingException} +import com.fasterxml.jackson.databind.exc.MismatchedInputException import org.apache.hadoop.fs.Path import org.scalatest.BeforeAndAfter @@ -30,6 +32,7 @@ import org.apache.spark._ import org.apache.spark.deploy.SparkHadoopUtil import org.apache.spark.deploy.history.EventLogFileReader import org.apache.spark.deploy.history.EventLogTestHelper._ +import org.apache.spark.internal.config.History import org.apache.spark.io.{CompressionCodec, LZ4CompressionCodec} import org.apache.spark.util.{JsonProtocol, JsonProtocolSuite, Utils} @@ -114,36 +117,62 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp val json = JsonProtocol.sparkEventToJsonString(start) assert(json.getBytes(StandardCharsets.UTF_8).length < 8 * 1024) val listener = new EventBufferingListener - val bus = new ReplayListenerBus(8 * 1024) + val conf = new SparkConf(false).set(History.EVENT_LOG_MAX_LINE_LENGTH.key, "8k") + val bus = new ReplayListenerBus(ReplayListenerBus.maxLineLength(conf)) bus.addListener(listener) val input = new ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8)) assert(bus.replay(input, "ascii")) assert(listener.loggedEvents.toSeq == Seq(json)) } - test("Replay line limit handles UTF-8 boundaries and line endings") { - val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + test("Replay line limit handles raw UTF-8 boundaries and line endings") { // scalastyle:off nonascii - val names = Seq("x", "\u00e9", "\u4e2d", "\ud83d\ude00") + val suffixes = Seq("x", "\u00e9", "\u4e2d", "\ud83d\ude00") // scalastyle:on nonascii for { - name <- names - repetitions <- Seq(100, 4097) + suffix <- suffixes + padding <- Seq(8, 8190, 8191, 8192) ending <- Seq("\n", "\r\n", "") - excess <- Seq(0, 1) + delta <- -3 to 2 } { - val start = SparkListenerApplicationStart(name * repetitions, None, 125L, "user", None) - val json = JsonProtocol.sparkEventToJsonString(start) - val limit = json.getBytes(StandardCharsets.UTF_8).length - excess - val listener = new EventBufferingListener + // Put the limit inside the final code point, and split UTF-8 sequences across read buffers. + val line = "x" * padding + suffix + val bytes = line.getBytes(StandardCharsets.UTF_8) + val limit = bytes.length + delta + val seen = new ArrayBuffer[String] val bus = new ReplayListenerBus(limit) - bus.addListener(listener) - // Put a short event first, including when the final line has no terminator. - val input = new ByteArrayInputStream((end + "\n" + json + ending) - .getBytes(StandardCharsets.UTF_8)) - assert(bus.replay(input, "utf8")) - val expected = if (excess == 0) Seq(end, json) else Seq(end) - assert(listener.loggedEvents.toSeq == expected, s"name=$name ending=$ending excess=$excess") + val input = new ByteArrayInputStream((line + ending).getBytes(StandardCharsets.UTF_8)) + withClue(s"codePoint=${suffix.codePointAt(0)} padding=$padding " + + s"ending=${ending.map(_.toInt).mkString(",")} delta=$delta: ") { + // Observe raw lines before JSON parsing, without letting Jackson escape the test input. + assert(bus.replay(input, "utf8", eventsFilter = line => { + seen += line + false + })) + assert(seen.toSeq == (if (delta >= 0) Seq(line) else Seq.empty)) + } + } + } + + test("Replay line buffer stops retaining oversized content") { + for (limit <- Seq(0, 1, 8, 1024)) { + val buffer = new ReplayListenerBus.BoundedLineBuffer(limit) + for (_ <- 0 until 10000) { + buffer.append('x') + assert(buffer.length <= limit + 1) + } + assert(buffer.length == limit + 1) + assert(buffer.result().isEmpty) + } + } + + test("Replay line limit defaults and maximum supported configuration") { + val conf = new SparkConf(false) + assert(conf.get(History.EVENT_LOG_MAX_LINE_LENGTH) == 256L * 1024 * 1024) + assert(ReplayListenerBus.maxLineLength(conf) == ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) + for (value <- Seq("0", "-1", "2147483647", "3g")) { + conf.set(History.EVENT_LOG_MAX_LINE_LENGTH.key, value) + assert(ReplayListenerBus.maxLineLength(conf) == Int.MaxValue) } } @@ -161,6 +190,44 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp .contains("Malformed line #4: {}"))) } + test("Replay logs physical line numbers before rethrowing JSON errors") { + val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + val mappingError = + """{"Event":"org.apache.spark.util.TestListenerEvent","foo":"x","bar":[]}""" + for ((malformed, errorClass) <- Seq( + ("{bad", classOf[JsonParseException]), + (mappingError, classOf[MismatchedInputException]))) { + val appender = new LogAppender + val bus = new ReplayListenerBus(1024) + val input = new ByteArrayInputStream( + Seq(end, "x" * 2048, "x" * 2048, malformed, end).mkString("\n") + .getBytes(StandardCharsets.UTF_8)) + withLogAppender(appender) { + val error = intercept[JsonProcessingException] { + bus.replay(input, "json-errors", maybeTruncated = true) + } + assert(errorClass.isInstance(error)) + } + assert(appender.loggingEvents.exists(_.getMessage.getFormattedMessage + .contains("Exception parsing Spark event log: json-errors at line 4"))) + } + } + + test("Replay propagates stream IO errors without reporting a stale parse line") { + val error = new IOException("read failure") + val input = new InputStream { + override def read(): Int = throw error + } + val appender = new LogAppender + withLogAppender(appender) { + assert(intercept[IOException] { + new ReplayListenerBus().replay(input, "read-error") + } eq error) + } + assert(!appender.loggingEvents.exists(_.getMessage.getFormattedMessage + .contains("Exception parsing Spark event log"))) + } + test("Replay preserves line numbers for truncated logs and filtered iterators") { val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) val truncated = "{\"Event\":" diff --git a/docs/monitoring.md b/docs/monitoring.md index c81b6569cfa09..b1d94e4180104 100644 --- a/docs/monitoring.md +++ b/docs/monitoring.md @@ -440,12 +440,13 @@ Security options for the Spark History Server are covered more detail in the spark.history.fs.eventLog.maxLineLength - 512m + 256m Maximum UTF-8 byte length of a single event log line during replay, excluding the line ending. Longer lines are skipped with a warning, which bounds 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.
+ 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. Values at or below + 0, or above 2147483647, use the maximum supported limit of 2147483647 bytes.
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. From aaee04effe9b46d244faaab1a7b6c39b89bbcd57 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 27 Sep 2026 20:05:47 -0700 Subject: [PATCH 3/5] [SPARK-59804][UI] Clarify optional replay line naming --- .../org/apache/spark/scheduler/ReplayListenerBus.scala | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala index db8984f41e546..6052b427b0d3a 100644 --- a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala +++ b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala @@ -119,8 +119,8 @@ private[spark] class ReplayListenerBus( buffer.append(c.toChar) c = reader.read() } - val line = buffer.result() - if (line.isEmpty) { + val maybeLine = buffer.result() + if (maybeLine.isEmpty) { if (!warned) { logWarning(log"Skipped event log lines longer than " + log"${MDC(MAX_SIZE, maxLineLength)} bytes in " + @@ -129,7 +129,7 @@ private[spark] class ReplayListenerBus( } fetchLine() } else { - (line.get, index) + (maybeLine.get, index) } } } From 691a1d7e6db3068a34c186532a86574357a9342a Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Mon, 28 Sep 2026 00:15:11 -0700 Subject: [PATCH 4/5] [SPARK-59804][UI] Cap replay line limits and avoid duplicate stack traces --- .../spark/internal/config/History.scala | 5 ++-- .../spark/scheduler/ReplayListenerBus.scala | 28 +++++++++++-------- .../spark/scheduler/ReplayListenerSuite.scala | 25 ++++++++++++++--- docs/monitoring.md | 3 +- 4 files changed, 43 insertions(+), 18 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/internal/config/History.scala b/core/src/main/scala/org/apache/spark/internal/config/History.scala index 83a28429fa79b..ab517ca2bf966 100644 --- a/core/src/main/scala/org/apache/spark/internal/config/History.scala +++ b/core/src/main/scala/org/apache/spark/internal/config/History.scala @@ -171,8 +171,9 @@ private[spark] object History { "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. Values at or below 0, or above 2147483647, use the maximum supported limit of " + - "2147483647 bytes. " + + "space. Values at or below 0, or above 536870912, use the maximum supported limit of " + + "536870912 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") diff --git a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala index 6052b427b0d3a..12d576a0a2cd1 100644 --- a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala +++ b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala @@ -43,6 +43,9 @@ private[spark] class ReplayListenerBus( maxLineLength: Int = ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) extends SparkListenerBus with Logging { + 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. @@ -108,7 +111,7 @@ private[spark] class ReplayListenerBus( } @tailrec private def fetchLine(): (String, Int) = { - val buffer = new BoundedLineBuffer(maxLineLength) + val buffer = new BoundedLineBuffer(effectiveMaxLineLength) var c = reader.read() if (c == -1) { null @@ -123,7 +126,7 @@ private[spark] class ReplayListenerBus( if (maybeLine.isEmpty) { if (!warned) { logWarning(log"Skipped event log lines longer than " + - log"${MDC(MAX_SIZE, maxLineLength)} bytes in " + + log"${MDC(MAX_SIZE, effectiveMaxLineLength)} bytes in " + log"${MDC(FILE_NAME, sourceName)}") warned = true } @@ -207,7 +210,7 @@ private[spark] class ReplayListenerBus( case _: EOFException if maybeTruncated => false case jpe: JsonProcessingException => logError(log"Exception parsing Spark event log: ${MDC(PATH, sourceName)} " + - log"at line ${MDC(LINE_NUM, lineNumber)}", jpe) + log"at line ${MDC(LINE_NUM, lineNumber)}") throw jpe case ioe: IOException => throw ioe @@ -232,16 +235,19 @@ private[spark] class HaltReplayException extends RuntimeException private[spark] object ReplayListenerBus { - /** - * Default UTF-8 line-length cap, matching spark.history.fs.eventLog.maxLineLength. - * Keep the maximum buffered character count close to the original 512 MiB / 2 cap. - */ - val DEFAULT_MAX_LINE_LENGTH: Int = 256 * 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 - /** Resolves the byte limit, using Int.MaxValue for non-positive or larger values. */ + // Bound StringBuilder growth so UTF-16 inflation stays below the JVM array-size limit. + val MAX_LINE_LENGTH: Int = 512 * 1024 * 1024 + + /** 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 diff --git a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala index a11db118870c9..867c2ac8f0022 100644 --- a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala @@ -170,12 +170,27 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp val conf = new SparkConf(false) assert(conf.get(History.EVENT_LOG_MAX_LINE_LENGTH) == 256L * 1024 * 1024) assert(ReplayListenerBus.maxLineLength(conf) == ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) - for (value <- Seq("0", "-1", "2147483647", "3g")) { + for (value <- Seq("0", "-1", "512m", "536870913", "576m", "600m", "1g", + "2147483647", "3g")) { conf.set(History.EVENT_LOG_MAX_LINE_LENGTH.key, value) - assert(ReplayListenerBus.maxLineLength(conf) == Int.MaxValue) + assert(ReplayListenerBus.maxLineLength(conf) == ReplayListenerBus.MAX_LINE_LENGTH) } } + test("Replay constructor normalizes the effective line limit") { + val maximum = ReplayListenerBus.MAX_LINE_LENGTH + for (limit <- Seq(Int.MinValue, -1, 0, maximum, maximum + 1, Int.MaxValue)) { + assert(new ReplayListenerBus(limit).effectiveMaxLineLength == maximum) + } + for (limit <- Seq(1, 8192, ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH, maximum - 1)) { + assert(new ReplayListenerBus(limit).effectiveMaxLineLength == limit) + val conf = new SparkConf(false).set(History.EVENT_LOG_MAX_LINE_LENGTH, limit.toLong) + assert(ReplayListenerBus.maxLineLength(conf) == limit) + } + assert(new ReplayListenerBus().effectiveMaxLineLength == + ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) + } + test("Replay preserves physical line numbers after skipping long lines") { val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) val appender = new LogAppender @@ -208,8 +223,10 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp } assert(errorClass.isInstance(error)) } - assert(appender.loggingEvents.exists(_.getMessage.getFormattedMessage - .contains("Exception parsing Spark event log: json-errors at line 4"))) + val diagnostics = appender.loggingEvents.filter(_.getMessage.getFormattedMessage + .contains("Exception parsing Spark event log: json-errors at line 4")) + assert(diagnostics.size == 1) + assert(diagnostics.head.getThrown == null) } } diff --git a/docs/monitoring.md b/docs/monitoring.md index b1d94e4180104..3bc73277525b6 100644 --- a/docs/monitoring.md +++ b/docs/monitoring.md @@ -446,7 +446,8 @@ Security options for the Spark History Server are covered more detail in the ending. Longer lines are skipped with a warning, which bounds 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. Values at or below - 0, or above 2147483647, use the maximum supported limit of 2147483647 bytes.
+ 0, or above 536870912, use the maximum supported limit of 536870912 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. From 1167bfad5de6b7da1fe634b4e268de0e041b1887 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Mon, 28 Sep 2026 13:22:03 -0700 Subject: [PATCH 5/5] [SPARK-59804][UI] Restore replay constructor and release skipped line buffers --- .../spark/internal/config/History.scala | 10 ++- .../spark/scheduler/ReplayListenerBus.scala | 20 ++++- .../spark/scheduler/ReplayListenerSuite.scala | 86 +++++++++++++++---- docs/monitoring.md | 4 +- 4 files changed, 98 insertions(+), 22 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/internal/config/History.scala b/core/src/main/scala/org/apache/spark/internal/config/History.scala index ab517ca2bf966..b75500f7e5c34 100644 --- a/core/src/main/scala/org/apache/spark/internal/config/History.scala +++ b/core/src/main/scala/org/apache/spark/internal/config/History.scala @@ -165,15 +165,19 @@ 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 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. Values at or below 0, or above 536870912, use the maximum supported limit of " + - "536870912 bytes (512 MiB). This cap avoids JVM array-size limits but does not " + - "guarantee sufficient heap space. " + + "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") diff --git a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala index 12d576a0a2cd1..e2d7892428f60 100644 --- a/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala +++ b/core/src/main/scala/org/apache/spark/scheduler/ReplayListenerBus.scala @@ -38,11 +38,15 @@ import org.apache.spark.util.JsonProtocol * @param maxLineLength Maximum UTF-8 byte length of a single event log line, excluding its line * 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( 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) @@ -127,9 +131,10 @@ private[spark] class ReplayListenerBus( if (!warned) { logWarning(log"Skipped event log lines longer than " + log"${MDC(MAX_SIZE, effectiveMaxLineLength)} bytes in " + - log"${MDC(FILE_NAME, sourceName)}") + 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 { (maybeLine.get, index) @@ -216,7 +221,10 @@ private[spark] class ReplayListenerBus( throw ioe case e: Exception => 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 } } @@ -239,7 +247,7 @@ private[spark] object ReplayListenerBus { 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 = 512 * 1024 * 1024 + val MAX_LINE_LENGTH: Int = History.EVENT_LOG_MAX_LINE_LENGTH_LIMIT /** Resolves the byte limit, using 512 MiB for non-positive or larger values. */ def maxLineLength(conf: SparkConf): Int = { @@ -263,6 +271,8 @@ private[spark] object ReplayListenerBus { 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. @@ -270,6 +280,10 @@ private[spark] object ReplayListenerBus { // Reserve one extra byte for a possible CR in the line ending. if (byteLength <= bufferLimit) { buffer.append(c) + } else { + // The line will be skipped; release the retained prefix while draining. + buffer.setLength(0) + buffer.trimToSize() } } } diff --git a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala index 867c2ac8f0022..f99a61eaa29f9 100644 --- a/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala +++ b/core/src/test/scala/org/apache/spark/scheduler/ReplayListenerSuite.scala @@ -26,6 +26,7 @@ import scala.collection.mutable.ArrayBuffer import com.fasterxml.jackson.core.{JsonParseException, JsonProcessingException} import com.fasterxml.jackson.databind.exc.MismatchedInputException import org.apache.hadoop.fs.Path +import org.apache.logging.log4j.Level import org.scalatest.BeforeAndAfter import org.apache.spark._ @@ -112,7 +113,7 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp assert(eventMonster.loggedEvents(1) === JsonProtocol.sparkEventToJsonString(applicationEnd)) } - test("Replay line limit uses UTF-8 bytes") { + test("SPARK-59804: Replay line limit uses UTF-8 bytes") { val start = SparkListenerApplicationStart("x" * 6000, None, 125L, "user", None) val json = JsonProtocol.sparkEventToJsonString(start) assert(json.getBytes(StandardCharsets.UTF_8).length < 8 * 1024) @@ -125,7 +126,7 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp assert(listener.loggedEvents.toSeq == Seq(json)) } - test("Replay line limit handles raw UTF-8 boundaries and line endings") { + test("SPARK-59804: Replay line limit handles raw UTF-8 boundaries and line endings") { // scalastyle:off nonascii val suffixes = Seq("x", "\u00e9", "\u4e2d", "\ud83d\ude00") // scalastyle:on nonascii @@ -154,19 +155,59 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp } } - test("Replay line buffer stops retaining oversized content") { + test("SPARK-59804: Replay line buffer stops retaining oversized content") { for (limit <- Seq(0, 1, 8, 1024)) { val buffer = new ReplayListenerBus.BoundedLineBuffer(limit) for (_ <- 0 until 10000) { buffer.append('x') assert(buffer.length <= limit + 1) } - assert(buffer.length == limit + 1) + assert(buffer.length == 0) + assert(buffer.capacity == 0) assert(buffer.result().isEmpty) } } - test("Replay line limit defaults and maximum supported configuration") { + test("SPARK-59804: Replay retains a public no-argument constructor") { + val bus = classOf[ReplayListenerBus].getConstructor().newInstance() + assert(bus.effectiveMaxLineLength == ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) + } + + test("SPARK-59804: Replay reports the physical location of skipped lines") { + val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + val input = new ByteArrayInputStream( + Seq(end, "x" * 2048, end, "x" * 2048).mkString("\n") + .getBytes(StandardCharsets.UTF_8)) + val appender = new LogAppender + appender.setThreshold(Level.DEBUG) + withLogAppender(appender, level = Some(Level.DEBUG)) { + assert(new ReplayListenerBus(1024).replay(input, "skipped-lines")) + } + val warnings = appender.loggingEvents.filter(_.getLevel == Level.WARN) + assert(warnings.size == 1) + assert(warnings.head.getMessage.getFormattedMessage.contains("first skipped line: 2")) + val debug = appender.loggingEvents.filter(_.getLevel == Level.DEBUG) + .map(_.getMessage.getFormattedMessage) + assert(debug.contains("Skipped event log line 2 in skipped-lines")) + assert(debug.contains("Skipped event log line 4 in skipped-lines")) + } + + test("SPARK-59804: Replay bounds malformed line diagnostics") { + val line = """{"Event":"SparkListenerJobStart","Job ID":1,"Stage IDs":[],""" + + """"Properties":{"bad":1},"padding":"""" + "x" * 10000 + "tail-marker\"}" + val appender = new LogAppender + withLogAppender(appender) { + assert(!new ReplayListenerBus().replay( + new ByteArrayInputStream(line.getBytes(StandardCharsets.UTF_8)), "semantic-error")) + } + val diagnostic = appender.loggingEvents.map(_.getMessage.getFormattedMessage) + .find(_.startsWith("Malformed line #1:")).get + assert(diagnostic.contains(s"line length: ${line.length} characters")) + assert(!diagnostic.contains("tail-marker")) + assert(diagnostic.length < 1200) + } + + test("SPARK-59804: Replay line limit defaults and maximum supported configuration") { val conf = new SparkConf(false) assert(conf.get(History.EVENT_LOG_MAX_LINE_LENGTH) == 256L * 1024 * 1024) assert(ReplayListenerBus.maxLineLength(conf) == ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) @@ -177,10 +218,17 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp } } - test("Replay constructor normalizes the effective line limit") { + test("SPARK-59804: Replay constructor normalizes the effective line limit") { val maximum = ReplayListenerBus.MAX_LINE_LENGTH for (limit <- Seq(Int.MinValue, -1, 0, maximum, maximum + 1, Int.MaxValue)) { - assert(new ReplayListenerBus(limit).effectiveMaxLineLength == maximum) + val bus = new ReplayListenerBus(limit) + assert(bus.effectiveMaxLineLength == maximum) + val listener = new EventBufferingListener + bus.addListener(listener) + val json = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + assert(bus.replay(new ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8)), + "normalized")) + assert(listener.loggedEvents.toSeq == Seq(json)) } for (limit <- Seq(1, 8192, ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH, maximum - 1)) { assert(new ReplayListenerBus(limit).effectiveMaxLineLength == limit) @@ -191,7 +239,7 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp ReplayListenerBus.DEFAULT_MAX_LINE_LENGTH) } - test("Replay preserves physical line numbers after skipping long lines") { + test("SPARK-59804: Replay preserves physical line numbers after skipping long lines") { val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) val appender = new LogAppender val bus = new ReplayListenerBus(1024) @@ -205,7 +253,7 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp .contains("Malformed line #4: {}"))) } - test("Replay logs physical line numbers before rethrowing JSON errors") { + test("SPARK-59804: Replay logs physical line numbers before rethrowing JSON errors") { val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) val mappingError = """{"Event":"org.apache.spark.util.TestListenerEvent","foo":"x","bar":[]}""" @@ -230,22 +278,30 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp } } - test("Replay propagates stream IO errors without reporting a stale parse line") { + test("SPARK-59804: Replay propagates stream IO errors without reporting a stale parse line") { val error = new IOException("read failure") - val input = new InputStream { - override def read(): Int = throw error + val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) + val input = new ByteArrayInputStream((end + "\n").getBytes(StandardCharsets.UTF_8)) { + override def read(bytes: Array[Byte], offset: Int, length: Int): Int = { + if (available() == 0) throw error + super.read(bytes, offset, length) + } } + val bus = new ReplayListenerBus() + val listener = new EventBufferingListener + bus.addListener(listener) val appender = new LogAppender withLogAppender(appender) { assert(intercept[IOException] { - new ReplayListenerBus().replay(input, "read-error") + bus.replay(input, "read-error") } eq error) } + assert(listener.loggedEvents.toSeq == Seq(end)) assert(!appender.loggingEvents.exists(_.getMessage.getFormattedMessage .contains("Exception parsing Spark event log"))) } - test("Replay preserves line numbers for truncated logs and filtered iterators") { + test("SPARK-59804: Replay preserves line numbers for truncated logs and filtered iterators") { val end = JsonProtocol.sparkEventToJsonString(SparkListenerApplicationEnd(1000L)) val truncated = "{\"Event\":" val appender = new LogAppender @@ -264,7 +320,7 @@ class ReplayListenerSuite extends SparkFunSuite with BeforeAndAfter with LocalSp assert(warnings.forall(_.contains("at line 3,"))) } - test("Replay still rejects malformed UTF-8 in skipped lines") { + test("SPARK-59804: Replay still rejects malformed UTF-8 in skipped lines") { val bytes = Array.fill[Byte](20 * 1024)('x'.toByte) ++ Array(0xff.toByte, '\n'.toByte) val bus = new ReplayListenerBus(1024) intercept[java.nio.charset.MalformedInputException] { diff --git a/docs/monitoring.md b/docs/monitoring.md index 3bc73277525b6..55cf5d7785093 100644 --- a/docs/monitoring.md +++ b/docs/monitoring.md @@ -445,7 +445,9 @@ Security options for the Spark History Server are covered more detail in the Maximum UTF-8 byte length of a single event log line during replay, excluding the line ending. Longer lines are skipped with a warning, which bounds 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. Values at or below + 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, or above 536870912, use the maximum supported limit of 536870912 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