Skip to content

Commit c5ce330

Browse files
JoshRosenLuciferYang
authored andcommitted
[SPARK-58507][CORE] Avoid redundant closure-class parsing in ClosureCleaner's indylambda path
### What changes were proposed in this pull request? Two changes to `ClosureCleaner`'s indylambda path: 1. Hoist the `getCapturedArgCount == 0` check above `Class.forName` and the ASM parse. It is an O(1) read of a `SerializedLambda` the method already holds, and when it is zero there is nothing to clean, so the expensive work was being done only to be discarded. Skipping the return-statement fail-fast for such closures is safe: a non-local return compiles to `throw new NonLocalReturnControl(key, value)` where `key` is allocated in the enclosing method, so a closure containing one necessarily captures that key. The null-`getCapturedArg(0)` bail-out, by contrast, deliberately stays _below_ the return-statement check: a closure with a non-local return can capture a null value as its first captured argument, so hoisting that bail-out too would silently skip the fail-fast (see the new regression test, which fails in that configuration). 2. Memoize the return-statement check per class in a `java.lang.ClassValue`. `ReturnStatementFinder` is replaced by a collecting variant (`ReturnStatementCollector`) so one parse answers for every method on the class, preserving the existing `$adapted` matching rule exactly. Scanning granularity is unchanged: a classfile cannot be parsed per-method (ASM's `accept` walks the whole file), and both the old and new visitors decode instructions only for `apply`/`$anonfun$` methods. The previous code paid that full-class walk on every `clean()` call; the memoized version pays it once per class. The only per-parse work lost is the old throw-on-first-match early abort, which fired only when the closure was about to be rejected anyway. The stored result is a set of offending method names, or the shared empty-set singleton for any class with no non-local returns. Why `ClassValue` rather than a map keyed on `Class`? A `ConcurrentHashMap[Class[_], _]` would strongly reference its keys and prevent classloader unloading, which matters here because Spark's REPL and `ChildFirstURLClassLoader` discard loaders routinely. `ClassValue` associates the state with the `Class` itself, so an entry becomes unreachable precisely when its class does; no unloading is blocked. Verified empirically: a populated entry does not prevent class or classloader collection, and 20000 sequentially created classes each given an entry left zero alive after GC with metaspace flat. `computeValue` may be invoked concurrently and more than once per class, which is safe because return-statement detection is a pure function of the class bytes. As a side effect this removes a latent NPE: `getClassReader` has always been allowed to return `null` (bytecode is not resource-accessible for e.g. LambdaMetafactory-generated classes, which is why SPARK-14540 added the same guard in `getInnerClosureClasses`), and the two return-statement call sites dereferenced it unconditionally. At these sites the argument is a capturing class, which in practice has readable bytecode, so the NPE is not known to fire; the memoized check now honors the contract and skips such classes with a debug log. The fail-fast is best-effort, so skipping cannot cause incorrect execution, only a later, less friendly error if the closure really contains a non-local return. ### Why are the changes needed? `RDD.collect()` passes `(iter: Iterator[T]) => iter.toArray` to `runJob`, which captures nothing, yet `SparkContext.runJob` cleans it unconditionally -- loading a class, reading a class file out of a JAR and running a full ASM parse, on every collect. Profiling a Spark test JVM (`SQLQueryTestSuite`) showed: * `ClosureCleaner` on the stack for 13.2% of CPU samples and 25.2% of all allocation samples; * 3251 `getClassReader` invocations over 21 distinct classes -- a 155x repeat ratio, because the lambdas passed to `runJob` are declared by a handful of Spark's own classes (`WholeStageCodegenExec`, `SparkContext`, `RDD`, `Dataset`), so every job re-parses the same bytecode; * 14.7% of indylambda cleans are non-capturing. The cost is a per-job overhead, measured at 9.8-13.2% of test-JVM CPU across six suites including `core/RDDSuite`, which involves no SQL at all. Note that the ClosureCleaner's indylambda path only modifies closures whose first captured argument is a Scala REPL line object or an instance of a class defined in an Ammonite session; for all other code it is purely a validator. Since modification is gated on a captured argument existing, a closure with `getCapturedArgCount == 0` can never be modified, so the hoist cannot remove cleaning, and the legacy non-indylambda path's cleaning is untouched by this PR (only its validator is memoized). ### Does this PR introduce _any_ user-facing change? No. Behaviour is unchanged, including the fail-fast `ReturnStatementInClosureException` for every capturing closure. ### How was this patch tested? * New regression test "return statements in closures capturing a null value are identified at cleaning time": its closure's first captured argument is a null local that precedes the `NonLocalReturnControl` key in the capture order (verified via javap: `$anonfun$run$16(String, Object, int)`), pinning the requirement that the capture-count hoist must not skip the return-statement check for capturing closures. It fails if the null-capture bail-out is hoisted above the check. * New unit test "hasReturnStatement identifies non-local returns per method": exercises the any-method query, the exact impl-method match, the `$adapted`-suffix resolution rule (scala/scala-dev#109), a non-matching method name, and a class with no non-local returns. * `ClosureCleanerSuite` + `ClosureCleanerSuite2`: 18/18 (`ClosureCleanerSuite`, 12/12, rerun locally against the final patch). * `ReplSuite` + `SingletonReplSuite`: 36/36, and `AmmoniteReplE2ESuite`: 1/1. These are the two paths where cleaning actually modifies closures: the indylambda path only rewrites Scala REPL and Ammonite closures, so every other suite exercises the branch where a regression would be invisible. * `DataFrameSuite`: 176/176. * Effect confirmed by profiling the patched build on the same suite: `xbean.asm9` falls from 11.02% of samples to 0.17%, `getClassReader` to zero. * The above tests / profiling was on a Linux VPS box. Independently reproduced on a second machine (macOS, JDK 21, JFR `profile` settings), comparing baseline and patched `ClosureCleaner` swapped in front of an otherwise identical classpath: - `core/RDDSuite` end-to-end (79/79 green in both): cleaning-related frames fell from 8.1% of execution samples to 0.0%, `getClassReader` calls from 1306 over 11 distinct classes (a 119x repeat ratio) to 13, and suite wall clock dropped ~8.5%. - A loop of minimal `collect()`/`count()` jobs showed where the time goes: the driver re-parses `SparkContext.class` (204 KB) and `RDD.class` (189 KB) about five times per job; per-iteration wall time halved (11.1 ms -> 5.5 ms). - Microbenchmark of `SparkClosureCleaner.clean()` alone (baseline cost scales with the capturing class's file size; the memoized path is size-independent): for a 9 KB capturing class, 46.4 us -> 2.0 us per call; for a 113 KB one, 340 us -> 4.8 us. Non-capturing: 44.9 us -> 0.28 us. The memo's worst case (a never-repeated closure class) pays the same single parse as before plus ~360 bytes; measured repeat ratios on real suites were 119x (`RDDSuite`) and 155x (`DataFrameSuite`). ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Opus 5; refined with Claude Fable 5 Closes #57710 from JoshRosen/SPARK-58507-closurecleaner-redundant-parsing. Authored-by: Josh Rosen <rosenville@gmail.com> Signed-off-by: yangjie01 <yangjie01@baidu.com>
1 parent 62c57b4 commit c5ce330

2 files changed

Lines changed: 119 additions & 22 deletions

File tree

‎common/utils/src/main/scala/org/apache/spark/util/ClosureCleaner.scala‎

Lines changed: 56 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import java.io.{ByteArrayInputStream, ByteArrayOutputStream}
2121
import java.lang.invoke.{LambdaMetafactory, MethodHandle, MethodHandleInfo, MethodHandles, MethodType, SerializedLambda}
2222
import java.lang.reflect.{Field, Modifier}
2323

24+
import scala.collection.immutable
2425
import scala.collection.mutable.{ArrayBuffer, Map, Queue, Set, Stack}
2526
import scala.jdk.CollectionConverters._
2627

@@ -34,6 +35,40 @@ import org.apache.spark.internal.Logging
3435
* A cleaner that renders closures serializable if they can be done so safely.
3536
*/
3637
private[spark] object ClosureCleaner extends Logging {
38+
/**
39+
* Per-class memo of which closure methods contain a non-local return, i.e. allocate a
40+
* `scala/runtime/NonLocalReturnControl`. The verdict is a pure function of the class's
41+
* immutable bytecode, so one ASM parse per class answers for every `clean()` call.
42+
*/
43+
private val methodsWithNonLocalReturn = new ClassValue[immutable.Set[String]] {
44+
override def computeValue(cls: Class[_]): immutable.Set[String] = {
45+
val collector = new ReturnStatementCollector
46+
val reader = getClassReader(cls)
47+
if (reader != null) {
48+
reader.accept(collector, 0)
49+
} else {
50+
logDebug(s"Cannot get class bytes for ${cls.getName}; skipping return-statement check")
51+
}
52+
collector.found.toSet
53+
}
54+
}
55+
56+
/** Whether `implMethodName` (any closure method, if `None`) of `cls` has a non-local return. */
57+
private[util] def hasReturnStatement(cls: Class[_], implMethodName: Option[String]): Boolean = {
58+
val found = methodsWithNonLocalReturn.get(cls)
59+
implMethodName match {
60+
case None => found.nonEmpty
61+
case Some(target) =>
62+
// Some lambdas get an "$adapted" boxing bridge (e.g. { _: Int => return; Seq() }) while
63+
// others do not ({ _: Int => return; true }, fully specialized). When the bridge
64+
// exists, the SerializedLambda's impl method is the bridge, which only delegates: the
65+
// `new NonLocalReturnControl` instruction is emitted in the underlying unadapted
66+
// method, so also match with the suffix stripped.
67+
// See https://github.com/scala/scala-dev/issues/109.
68+
found.contains(target) || found.contains(target.stripSuffix("$adapted"))
69+
}
70+
}
71+
3772
// Get an ASM class reader for a given class from the JAR that loaded it
3873
private[util] def getClassReader(cls: Class[_]): ClassReader = {
3974
// Copy data over, before delegating to ClassReader - else we can run out of open file handles.
@@ -226,6 +261,13 @@ private[spark] object ClosureCleaner extends Logging {
226261

227262
logDebug(s"Cleaning indylambda closure: $implMethodName")
228263

264+
// A closure with no captured arguments needs no cleaning, and cannot contain a non-local
265+
// return (the `NonLocalReturnControl` key is allocated in the enclosing method and would
266+
// be captured), so return before loading and parsing the capturing class.
267+
if (lambdaProxy.getCapturedArgCount == 0) {
268+
return None
269+
}
270+
229271
// capturing class is the class that declared this lambda
230272
val capturingClassName = lambdaProxy.getCapturingClass.replace('/', '.')
231273
val classLoader = func.getClass.getClassLoader // this is the safest option
@@ -234,16 +276,13 @@ private[spark] object ClosureCleaner extends Logging {
234276
// scalastyle:on classforname
235277

236278
// Fail fast if we detect return statements in closures
237-
val capturingClassReader = getClassReader(capturingClass)
238-
capturingClassReader.accept(new ReturnStatementFinder(Option(implMethodName)), 0)
239-
240-
val outerThis = if (lambdaProxy.getCapturedArgCount > 0) {
241-
// only need to clean when there is an enclosing non-null "this" captured by the closure
242-
Option(lambdaProxy.getCapturedArg(0)).getOrElse(return None)
243-
} else {
244-
return None
279+
if (hasReturnStatement(capturingClass, Option(implMethodName))) {
280+
throw new ReturnStatementInClosureException
245281
}
246282

283+
// only need to clean when there is an enclosing non-null "this" captured by the closure
284+
val outerThis = Option(lambdaProxy.getCapturedArg(0)).getOrElse(return None)
285+
247286
// clean only if enclosing "this" is something cleanable, i.e. a Scala REPL line object or
248287
// Ammonite command helper object.
249288
// For Ammonite closures, we do not care about actual capturing class name,
@@ -312,7 +351,9 @@ private[spark] object ClosureCleaner extends Logging {
312351
}
313352

314353
// Fail fast if we detect return statements in closures
315-
getClassReader(func.getClass).accept(new ReturnStatementFinder(), 0)
354+
if (hasReturnStatement(func.getClass, None)) {
355+
throw new ReturnStatementInClosureException
356+
}
316357

317358
// If accessed fields is not populated yet, we assume that
318359
// the closure we are trying to clean is the starting one
@@ -1075,26 +1116,19 @@ private[spark] object IndylambdaScalaClosures extends Logging {
10751116
private[spark] class ReturnStatementInClosureException
10761117
extends SparkException("Return statements aren't allowed in Spark closures")
10771118

1078-
private class ReturnStatementFinder(targetMethodName: Option[String] = None)
1079-
extends ClassVisitor(Opcodes.ASM9) {
1119+
/** Collects the names of all closure methods that contain a non-local return. */
1120+
private class ReturnStatementCollector extends ClassVisitor(Opcodes.ASM9) {
1121+
val found = Set.empty[String]
1122+
10801123
override def visitMethod(access: Int, name: String, desc: String,
10811124
sig: String, exceptions: Array[String]): MethodVisitor = {
10821125

10831126
// $anonfun$ covers indylambda closures
10841127
if (name.contains("apply") || name.contains("$anonfun$")) {
1085-
// A method with suffix "$adapted" will be generated in cases like
1086-
// { _:Int => return; Seq()} but not { _:Int => return; true}
1087-
// closure passed is $anonfun$t$1$adapted while actual code resides in $anonfun$s$1
1088-
// visitor will see only $anonfun$s$1$adapted, so we remove the suffix, see
1089-
// https://github.com/scala/scala-dev/issues/109
1090-
val isTargetMethod = targetMethodName.isEmpty ||
1091-
name == targetMethodName.get || name == targetMethodName.get.stripSuffix("$adapted")
1092-
10931128
new MethodVisitor(Opcodes.ASM9) {
10941129
override def visitTypeInsn(op: Int, tp: String): Unit = {
1095-
if (op == Opcodes.NEW && tp.contains("scala/runtime/NonLocalReturnControl") &&
1096-
isTargetMethod) {
1097-
throw new ReturnStatementInClosureException
1130+
if (op == Opcodes.NEW && tp.contains("scala/runtime/NonLocalReturnControl")) {
1131+
found += name
10981132
}
10991133
}
11001134
}

‎core/src/test/scala/org/apache/spark/util/ClosureCleanerSuite.scala‎

Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,11 +60,40 @@ class ClosureCleanerSuite extends SparkFunSuite {
6060
}
6161
}
6262

63+
test("return statements in closures capturing a null value are identified at cleaning time") {
64+
intercept[ReturnStatementInClosureException] {
65+
TestObjectWithBogusReturnsAndNullCapture.run()
66+
}
67+
}
68+
6369
test("return statements from named functions nested in closures don't raise exceptions") {
6470
val result = TestObjectWithNestedReturns.run()
6571
assert(result === 1)
6672
}
6773

74+
test("hasReturnStatement identifies non-local returns per method") {
75+
TestObjectWithReturnInClosure.run()
76+
val cls = TestObjectWithReturnInClosure.getClass
77+
val implMethodName = {
78+
val proxy =
79+
IndylambdaScalaClosures.getSerializationProxy(TestObjectWithReturnInClosure.lastClosure)
80+
assert(proxy.isDefined)
81+
proxy.get.getImplMethodName
82+
}
83+
// Any-method query.
84+
assert(ClosureCleaner.hasReturnStatement(cls, None))
85+
// Targeted query with the exact impl method name.
86+
assert(ClosureCleaner.hasReturnStatement(cls, Some(implMethodName)))
87+
// An "$adapted" wrapper name resolves to the underlying method that holds the closure body
88+
// (see https://github.com/scala/scala-dev/issues/109).
89+
assert(ClosureCleaner.hasReturnStatement(cls, Some(implMethodName + "$adapted")))
90+
// A method that does not exist on the class must not match.
91+
assert(!ClosureCleaner.hasReturnStatement(cls, Some("$anonfun$doesNotExist$1")))
92+
// A class whose closures contain no non-local returns.
93+
TestObjectWithoutReturnInClosure.run()
94+
assert(!ClosureCleaner.hasReturnStatement(TestObjectWithoutReturnInClosure.getClass, None))
95+
}
96+
6897
test("user provided closures are actually cleaned") {
6998

7099
// We use return statements as an indication that a closure is actually being cleaned
@@ -199,6 +228,40 @@ object TestObjectWithBogusReturns {
199228
}
200229
}
201230

231+
object TestObjectWithBogusReturnsAndNullCapture {
232+
def run(): Int = {
233+
withSpark(new SparkContext("local", "test")) { sc =>
234+
val nums = sc.parallelize(Array(1, 2, 3, 4).toImmutableArraySeq)
235+
val s: String = null
236+
// The closure's first captured argument may be the (null) `s` rather than the non-local
237+
// return's key: the cleaner must still detect the invalid return rather than bail out on
238+
// the null capture.
239+
nums.map { x => if (s != null) return 1; x * 2 }
240+
1
241+
}
242+
}
243+
}
244+
245+
object TestObjectWithReturnInClosure {
246+
// The non-local `return` forces the enclosing method to have type Int, so it cannot return the
247+
// closure to the caller directly; stash it for the test to inspect instead.
248+
var lastClosure: Int => Int = null
249+
def run(): Int = {
250+
val f = (x: Int) => { if (x < 0) return -1; x }
251+
lastClosure = f
252+
f(1)
253+
}
254+
}
255+
256+
object TestObjectWithoutReturnInClosure {
257+
var lastClosure: Int => Int = null
258+
def run(): Int = {
259+
val f = (x: Int) => x * 2
260+
lastClosure = f
261+
f(1)
262+
}
263+
}
264+
202265
object TestObjectWithNestedReturns {
203266
def run(): Int = {
204267
withSpark(new SparkContext("local", "test")) { sc =>

0 commit comments

Comments
 (0)