apache / apache/auron

ClassCastException (Scala 2.13 DefaultSerializationProxy) when deserializing ScalarSubquery injected by runtime bloom filter on Spark 4.1

Open
#2,386 3 comments 0 reactions 1 assignee Claimed by @xiaoyanxie View on GitHub
Dominant language
Rust
Stars
1.8k
Forks
241
Avg merge
2d 14h
Merged PRs (30d)
22

Description

**Describe the bug**

On Spark 4.1 (Scala 2.13), queries fail at runtime whenever Spark's runtime bloom filter optimization (`InjectRuntimeFilter`) is applied. Every affected task dies with:

```
java.lang.ClassCastException: cannot assign instance of scala.collection.generic.DefaultSerializationProxy
to field org.apache.spark.rdd.RDD.dependencies_ of type scala.collection.immutable.Seq
in instance of org.apache.spark.rdd.MapPartitionsRDD
at java.base/java.io.ObjectStreamClass$FieldReflector.setObjFieldValues(ObjectStreamClass.java:2096)
...
at org.apache.spark.sql.auron.NativeConverters$.deserializeExpression(NativeConverters.scala:1495)
at org.apache.auron.spark.sql.SparkAuronUDFWrapperContext.(SparkAuronUDFWrapperContext.scala:44)
at org.apache.auron.jni.SparkAuronAdaptor.getAuronUDFWrapperContext(SparkAuronAdaptor.java:103)
at org.apache.auron.jni.JniBridge.getAuronUDFWrapperContext(JniBridge.java:113)
```

followed by a native panic (`output_with_sender[Filter]: output() returns error: External error: Java exception thrown at native-engine/datafusion-ext-exprs/src/spark_udf_wrapper.rs:97`) and job abort.

Root cause (verified, see repro below):

1. `InjectRuntimeFilter` rewrites join filters into `might_contain(scalar-subquery(bloom_agg), xxhash64(key))`. Auron converts `BloomFilterMightContain` natively (`ShimsImpl.convertBloomFilterMightContain`), and its `bloomFilterExpression` — an `ExecSubqueryExpression` — is converted at `NativeConverters.scala:462`, which Java-serializes **the entire execution-side `ScalarSubquery` object, including its `plan: BaseSubqueryExec` reference**. An executed plan transitively references RDDs (`MapPartitionsRDD` in the trace above).
2. On the native side, `SparkScalarSubqueryWrapperExpr` evaluates by delegating to `SparkUDFWrapperExpr` (`native-engine/datafusion-ext-exprs/src/spark_scalar_subquery_wrapper.rs`), which calls back into the JVM and deserializes the bytes with a plain `java.io.ObjectInputStream` (`NativeConverters.deserializeExpression`). This is why the stack shows `SparkAuronUDFWrapperContext` although no UDF fallback is involved.
3. On Scala 2.13, immutable collections serialize through `scala.collection.generic.DefaultSerializationProxy`, which is only replaced by the real collection in `readResolve()`. Per the Java serialization spec, `readResolve` substitution does not propagate to back-references, so in a cyclic graph (RDD lineage) the *unresolved proxy* gets assigned into `RDD.dependencies_` → `ClassCastException`. Scala 2.12 has no serialization proxies for these collections, which is why the identical graph deserializes fine on 2.12.

Note this is **not** the `SparkUDFWrapper` expression-fallback path: ordinary fallback expressions (e.g. `UnscaledValue` / `MakeDecimal`, 167 instances in these runs) serialize shallow bound expression trees with no plan references and round-trip fine on Scala 2.13.

**To Reproduce**

Minimal repro (small dataset, threshold lowered so the runtime bloom filter is injected):

1. Build Auron for Spark 4.1 / Scala 2.13 (`./auron-build.sh --pre --sparkver 4.1 --scalaver 2.13`).
2. Generate TPC-DS parquet data at sf=10 (~5 GB).
3. Run:
```
SPARK_HOME= SPARK_VERSION=spark-4.1 SCALA_VERSION=2.13 \
dev/auron-it/run-it.sh --type tpcds --data-location \
--conf spark.sql.optimizer.runtime.bloomFilter.applicationSideScanSizeThreshold=1GB \
--query-filter q1,q2,q3
```
4. q2 and q3 fail with the exception above (the Auron run produces 0 rows, reported as a result mismatch vs the vanilla Spark baseline). q1 passes — its scanned tables stay below even the lowered threshold, so no runtime filter is injected there.

With default Spark settings the bug reproduces at sf=100, where the relevant fact-table scans cross the default `spark.sql.optimizer.runtime.bloomFilter.applicationSideScanSizeThreshold` (10GB):

| per-table scan size (uncompressed parquet) | sf=30 | sf=60 | sf=100 |
|---|---|---|---|
| `store_sales` (q3 application side) | 4.7 GB | 9.5 GB | **16 GB** |
| `catalog_sales` (q2 application side) | 3.3 GB | 6.9 GB | **12 GB** |

Verification matrix (Spark 4.1.2 / Scala 2.13, same host, JDK 21; CCE = `ClassCastException` count in log):

| | sf=1 | sf=10 | sf=10, threshold=1GB | sf=30 | sf=60 | sf=100 | sf=100, bloomFilter.enabled=false |
|---|---|---|---|---|---|---|---|
| result | PASS 3/3 | PASS 3/3 | **FAIL q2,q3** | PASS 3/3 | PASS 3/3 | **FAIL q2,q3** | PASS 3/3 |
| CCE | 0 | 0 | 83 | 0 | 0 | 52 | 0 |

Spark 3.5.8 / Scala 2.12 passes 3/3 at every scale factor with default settings (same code path serializes the same plan-carrying subquery, but Scala 2.12 collections have no serialization proxies).

The two toggles isolate the cause in both directions: disabling the runtime bloom filter at sf=100 eliminates the crash with everything else unchanged (including all 167 fallback-expression serializations), and lowering the threshold at sf=10 introduces it.

Independently re-verified on a second machine/OS (sf=10, threshold=1GB case): macOS aarch64, JDK 17.0.16, Scala 2.13, Spark 4.1.2, Auron 8.0.0-incubating built from source (`./auron-build.sh --pre --sparkver 4.1 --scalaver 2.13`). Result: `q1 PASS, q2 FAIL, q3 FAIL`, 66 `ClassCastException` occurrences, 25 ordinary fallback-expression serializations (`UnscaledValue`/`MakeDecimal`) — same qualitative signature as the Linux/JDK21 run (CCE count differs, as expected, since partitioning/task counts differ across environments/runs). The captured stack matches frame-for-frame:

```
java.lang.ClassCastException: cannot assign instance of scala.collection.generic.DefaultSerializationProxy
to field org.apache.spark.rdd.RDD.dependencies_ of type scala.collection.immutable.Seq
in instance of org.apache.spark.rdd.MapPartitionsRDD
...
at org.apache.spark.sql.auron.NativeConverters$.read$1(NativeConverters.scala:1498)
at org.apache.spark.sql.auron.NativeConverters$.$anonfun$deserializeExpression$4(NativeConverters.scala:1506)
at org.apache.spark.sql.auron.NativeConverters$.deserializeExpression(NativeConverters.scala:1495)
at org.apache.auron.spark.sql.SparkAuronUDFWrapperContext.(SparkAuronUDFWrapperContext.scala:44)
at org.apache.auron.jni.SparkAuronAdaptor.getAuronUDFWrapperContext(SparkAuronAdaptor.java:103)
at org.apache.auron.jni.JniBridge.getAuronUDFWrapperContext(JniBridge.java:113)

thread 'auron-native-stage-141-part-0-tid-251' panicked at native-engine/datafusion-ext-plans/src/common/execution_context.rs:573:21:
output_with_sender[Filter]: output() returns error: External error: Java exception thrown at native-engine/datafusion-ext-exprs/src/spark_udf_wrapper.rs:97: java.lang.ClassCastException: ...
```

This rules out a JDK-version-specific or OS/arch-specific explanation: the bug reproduces identically on Linux x86_64/JDK21 and macOS aarch64/JDK17, consistent with the root cause being purely about the Scala 2.13 collections runtime (`DefaultSerializationProxy`), not the JVM or platform.

**Expected behavior**

Queries with runtime-bloom-filter-injected scalar subqueries should run correctly on Spark 4.1 / Scala 2.13, as they do on Spark 3.5 / Scala 2.12. TPC-DS q2/q3 should pass the consistency check at sf=100 with default Spark settings.

**Screenshots**

N/A (full logs available on request).

**Additional context**

- Environment: Spark 4.1.2, Scala 2.13, JDK 21 (GraalVM 21.0.8), Auron master (8.0.0-incubating), local mode, Linux x86_64.
- Suggested fix direction: `NativeConverters.prepareExecSubquery` already calls `updateResult()` before serialization, so the subquery result is materialized on the driver at that point. Serializing the *evaluated result* (e.g. as a `Literal`) instead of the `ScalarSubquery` object would avoid capturing the physical plan entirely — fixing this crash and also removing plan/RDD payloads from task binaries (related to #2351 / #2307). Alternatively/additionally, a fail-fast check that serialized expressions contain no plan-referencing nodes would turn future regressions of this class into clear errors instead of Scala-version-dependent `ClassCastException`s.
- Related issues:
- #1953 — Spark 4.1 support epic (this only bites on Scala 2.13, hence only Spark 4.x).
- #2351 / #2307 — oversized task serialization from over-capture; same underlying over-capture, different symptom.
- #369 / #648 — earlier `ClassCastException`s in the same `deserializeExpression` path (classloader-related), showing this boundary is sensitive to the serialization environment.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.