From 6ec347f9494861866a3419d70200f1b186643576 Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 15:34:37 +0000 Subject: [PATCH 1/8] [Spark Runner] Replace Scala/shaded-Guava calls deprecated in Spark 3, removed in Spark 4 - runners/spark/src/.../io/{SourceRDD,SparkUnboundedSource}.java: scala.collection.JavaConversions -> JavaConverters (5 call-sites). - runners/spark/src/.../stateful/SparkGroupAlsoByWindowViaWindowSet.java: same JavaConversions -> JavaConverters migration (3 call-sites). - runners/spark/src/.../translation/batch/DoFnRunnerFactory.java: import scala.Serializable -> import java.io.Serializable (Scala 2.13 deprecated scala.Serializable; both are marker interfaces). - runners/spark/src/.../translation/{SparkStreamingPortablePipelineTranslator, streaming/StreamingTransformTranslator}.java: StreamingContext.union(asScalaBuffer(dStreams)) -> StreamingContext.union(asScalaBuffer(dStreams).toList()), since the mutable-buffer overload of union was deprecated in Spark 3 and removed in Spark 4 in favor of the immutable.Seq overload. - runners/spark/src/.../translation/streaming/ParDoStateUpdateFn.java: org.sparkproject.guava.collect.Iterators.emptyIterator() -> java.util.Collections.emptyIterator() (Spark 4 removed the legacy shaded-Guava Iterators). - runners/spark/src/.../structuredstreaming/SparkStructuredStreamingPipelineResult.java: static import of org.sparkproject.guava.base.Objects.firstNonNull -> vendor.guava MoreObjects.firstNonNull (Spark 4 dropped the legacy shaded Guava Objects class). Behavior-identical on Spark 3.5 (every replacement is a Scala/Guava API that exists with identical semantics on both Spark versions). First split-out commit from #38255 per @Abacn's review guidance. Also adds a CHANGES.md entry under "New Features / Improvements". --- CHANGES.md | 6 ++++++ .../org/apache/beam/runners/spark/io/SourceRDD.java | 10 +++++----- .../beam/runners/spark/io/SparkUnboundedSource.java | 2 +- .../stateful/SparkGroupAlsoByWindowViaWindowSet.java | 6 +++--- .../SparkStructuredStreamingPipelineResult.java | 2 +- .../translation/batch/DoFnRunnerFactory.java | 2 +- .../SparkStreamingPortablePipelineTranslator.java | 3 ++- .../translation/streaming/ParDoStateUpdateFn.java | 4 ++-- .../streaming/StreamingTransformTranslator.java | 2 +- 9 files changed, 22 insertions(+), 15 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index e8e11e830d14..d80b6affbb92 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -76,6 +76,12 @@ ([#38139](https://github.com/apache/beam/issues/38139)). * (Python) Added type alias for with_exception_handling to be used for typehints. ([#38173](https://github.com/apache/beam/issues/38173)). * Added plugin mechanism to support different Lineage implementations (Java) ([#36790](https://github.com/apache/beam/issues/36790)). +* Prepared the shared Spark runner base for Spark 4 compatibility: migrated + Scala collection and shaded-Guava calls to forms valid on both Spark 3 and + Spark 4, introduced a numeric `isSparkAtLeast` Gradle helper to replace + lexicographic Spark version comparison, and routed `requireJavaVersion` + through `applyJavaNature` so future Spark 4 builds enforce Java 17. No + behavior change on Spark 3.5. ## Breaking Changes diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SourceRDD.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SourceRDD.java index e65dccd23f24..56a2219933b6 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SourceRDD.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SourceRDD.java @@ -50,7 +50,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import scala.Option; -import scala.collection.JavaConversions; +import scala.collection.JavaConverters; /** Classes implementing Beam {@link Source} {@link RDD}s. */ @SuppressWarnings({ @@ -75,7 +75,7 @@ public static class Bounded extends RDD> { // to satisfy Scala API. private static final scala.collection.immutable.Seq> NIL = - JavaConversions.asScalaBuffer(Collections.>emptyList()).toList(); + JavaConverters.asScalaBuffer(Collections.>emptyList()).toList(); public Bounded( SparkContext sc, @@ -148,7 +148,7 @@ public scala.collection.Iterator> compute( final Iterator> readerIterator = new ReaderToIteratorAdapter<>(metricsContainer, reader); - return new InterruptibleIterator<>(context, JavaConversions.asScalaIterator(readerIterator)); + return new InterruptibleIterator<>(context, JavaConverters.asScalaIterator(readerIterator)); } /** @@ -299,7 +299,7 @@ public static class Unbounded> NIL = - JavaConversions.asScalaBuffer(Collections.>emptyList()).toList(); + JavaConverters.asScalaBuffer(Collections.>emptyList()).toList(); public Unbounded( SparkContext sc, @@ -344,7 +344,7 @@ public scala.collection.Iterator, CheckpointMarkT>> compu (CheckpointableSourcePartition) split; scala.Tuple2, CheckpointMarkT> tuple2 = new scala.Tuple2<>(partition.getSource(), partition.checkpointMark); - return JavaConversions.asScalaIterator(Collections.singleton(tuple2).iterator()); + return JavaConverters.asScalaIterator(Collections.singleton(tuple2).iterator()); } } diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SparkUnboundedSource.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SparkUnboundedSource.java index bea1557a7103..3f1fb103e47c 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SparkUnboundedSource.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/io/SparkUnboundedSource.java @@ -186,7 +186,7 @@ public Duration slideDuration() { @Override public scala.collection.immutable.List> dependencies() { - return scala.collection.JavaConversions.asScalaBuffer( + return scala.collection.JavaConverters.asScalaBuffer( Collections.>singletonList(parent)) .toList(); } diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkGroupAlsoByWindowViaWindowSet.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkGroupAlsoByWindowViaWindowSet.java index 2c54f90badbe..1c7a4c2a2416 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkGroupAlsoByWindowViaWindowSet.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkGroupAlsoByWindowViaWindowSet.java @@ -73,7 +73,7 @@ import scala.Tuple2; import scala.Tuple3; import scala.collection.Iterator; -import scala.collection.JavaConversions; +import scala.collection.JavaConverters; import scala.collection.Seq; import scala.runtime.AbstractFunction1; @@ -238,7 +238,7 @@ private Collection filterTimersEligibleForProcessing( // new input for key. try { final Iterable> elements = - FluentIterable.from(JavaConversions.asJavaIterable(encodedElements)) + FluentIterable.from(JavaConverters.asJavaIterable(encodedElements)) .transform(bytes -> CoderHelpers.fromByteArray(bytes, wvCoder)); LOG.trace("{}: input elements: {}", logPrefix, elements); @@ -410,7 +410,7 @@ private Collection filterTimersEligibleForProcessing( droppedDueToClosedWindow.inc(-droppedDueToClosedWindow.getCumulative()); } - return scala.collection.JavaConversions.asScalaIterator( + return JavaConverters.asScalaIterator( new UpdateStateByKeyOutputIterator(input, reduceFn, droppedDueToLateness)); } } diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java index b490ff875c31..806d838d9bff 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java @@ -18,7 +18,7 @@ package org.apache.beam.runners.spark.structuredstreaming; import static org.apache.beam.runners.core.metrics.MetricsContainerStepMap.asAttemptedOnlyMetricResults; -import static org.sparkproject.guava.base.Objects.firstNonNull; +import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects.firstNonNull; import java.io.IOException; import java.util.concurrent.ExecutionException; diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerFactory.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerFactory.java index 15ec818dba74..99ce3dc69889 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerFactory.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerFactory.java @@ -17,6 +17,7 @@ */ package org.apache.beam.runners.spark.structuredstreaming.translation.batch; +import java.io.Serializable; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -49,7 +50,6 @@ import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists; import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Maps; import org.joda.time.Instant; -import scala.Serializable; /** * Factory to create a {@link DoFnRunner}. The factory supports fusing multiple {@link DoFnRunner diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingPortablePipelineTranslator.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingPortablePipelineTranslator.java index 4850f886241b..9975c81b56a4 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingPortablePipelineTranslator.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingPortablePipelineTranslator.java @@ -330,7 +330,8 @@ private static void translateFlatten( } } // Unify streams into a single stream. - unifiedStreams = context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams)); + unifiedStreams = + context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams).toList()); } context.pushDataset( diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/ParDoStateUpdateFn.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/ParDoStateUpdateFn.java index 909624c23239..ed9299db4ee7 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/ParDoStateUpdateFn.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/ParDoStateUpdateFn.java @@ -19,6 +19,7 @@ import java.io.Serializable; import java.util.Collection; +import java.util.Collections; import java.util.Iterator; import java.util.List; import java.util.Map; @@ -62,7 +63,6 @@ import org.checkerframework.checker.nullness.qual.Nullable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.sparkproject.guava.collect.Iterators; import scala.Option; import scala.Tuple2; import scala.runtime.AbstractFunction3; @@ -236,7 +236,7 @@ public TimerInternals timerInternals() { final byte[] byteValue = serializedValue.get(); @Nullable WindowedValue windowedValue; @Nullable WindowedValue> keyedWindowedValue; - Iterator>> iterator = Iterators.emptyIterator(); + Iterator>> iterator = Collections.emptyIterator(); if (byteValue.length > 0) { windowedValue = CoderHelpers.fromByteArray(byteValue, this.wvCoder); keyedWindowedValue = windowedValue.withValue(KV.of(key, windowedValue.getValue())); diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StreamingTransformTranslator.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StreamingTransformTranslator.java index 48697a3dbafc..4a96edceba31 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StreamingTransformTranslator.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StreamingTransformTranslator.java @@ -306,7 +306,7 @@ public void evaluate(Flatten.PCollections transform, EvaluationContext contex } // start by unifying streams into a single stream. JavaDStream> unifiedStreams = - context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams)); + context.getStreamingContext().union(JavaConverters.asScalaBuffer(dStreams).toList()); context.putDataset(transform, new UnboundedDataset<>(unifiedStreams, streamingSources)); } From fda0505995e5b2a54dbdcc3369606a39438abf73 Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 15:34:53 +0000 Subject: [PATCH 2/8] [Spark Runner] Add Spark 4 hooks to BeamModulePlugin Three additive entries in buildSrc BeamModulePlugin.groovy: - def spark4_version = "4.0.2" alongside spark2_version and spark3_version. - project.ext.spark4_version export so per-project gradle scripts can reference it. - jackson_module_scala_2_13 library entry alongside the existing _2.11 and _2.12 entries. The first two are inert until the Spark 4 module landing PR (#38255) adds runners/spark/4/build.gradle that consumes them; the jackson_module_scala_2_13 library entry is independently useful for any future Scala 2.13 module. Second split-out commit from #38255 per @Abacn's review guidance. --- .../main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy | 3 +++ 1 file changed, 3 insertions(+) diff --git a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy index 005a8b587804..a99e0de47934 100644 --- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy +++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy @@ -649,6 +649,7 @@ class BeamModulePlugin implements Plugin { def solace_version = "10.21.0" def spark2_version = "2.4.8" def spark3_version = "3.5.0" + def spark4_version = "4.0.2" def spotbugs_version = "4.8.3" def testcontainers_version = "1.21.4" // [bomupgrader] determined by: org.apache.arrow:arrow-memory-core, consistent with: google_cloud_platform_libraries_bom @@ -658,6 +659,7 @@ class BeamModulePlugin implements Plugin { // Export Spark versions, so they are defined in a single place only project.ext.spark3_version = spark3_version + project.ext.spark4_version = spark4_version // version for BigQueryMetastore catalog (used by sdks:java:io:iceberg:bqms) // TODO: remove this and download the jar normally when the catalog gets // open-sourced (https://github.com/apache/iceberg/pull/11039) @@ -820,6 +822,7 @@ class BeamModulePlugin implements Plugin { jackson_datatype_jsr310 : "com.fasterxml.jackson.datatype:jackson-datatype-jsr310:$jackson_version", jackson_module_scala_2_11 : "com.fasterxml.jackson.module:jackson-module-scala_2.11:$jackson_version", jackson_module_scala_2_12 : "com.fasterxml.jackson.module:jackson-module-scala_2.12:$jackson_version", + jackson_module_scala_2_13 : "com.fasterxml.jackson.module:jackson-module-scala_2.13:$jackson_version", jamm : 'com.github.jbellis:jamm:0.4.0', jaxb_api : "jakarta.xml.bind:jakarta.xml.bind-api:$jaxb_api_version", jaxb_impl : "com.sun.xml.bind:jaxb-impl:$jaxb_api_version", From 5843a033a7f341ea6d522890651de9c71218f07f Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 15:38:17 +0000 Subject: [PATCH 3/8] [Spark Runner] Numeric isSparkAtLeast version compare; route requireJavaVersion to Spark 4 builds runners/spark/spark_runner.gradle: - Pass requireJavaVersion: (spark_version.startsWith("4") ? JavaVersion.VERSION_17 : null) into applyJavaNature so future Spark 4 builds enforce Java 17. No-op for Spark 3.5 (returns null). - Introduce isSparkAtLeast(minVersion) closure that compares numerically, e.g. so "3.10.0" sorts after "3.5.0" instead of before it lexicographically. - Replace all 5 call-sites of `if ("$spark_version" >= "3.5.0")` with `if (isSparkAtLeast("3.5.0"))`. Pure refactor, identical behaviour on the currently-supported Spark 3.x range. runners/spark/job-server/spark_job_server.gradle: - Same requireJavaVersion: arg routing for the job-server, gated on the parent spark_version property. requireJavaVersion itself is pre-existing in BeamModulePlugin.groovy (JavaNatureConfiguration parameter); this commit only adds invocation. Third split-out commit from #38255 per @Abacn's review guidance. --- .../spark/job-server/spark_job_server.gradle | 3 +++ runners/spark/spark_runner.gradle | 21 ++++++++++++++----- 2 files changed, 19 insertions(+), 5 deletions(-) diff --git a/runners/spark/job-server/spark_job_server.gradle b/runners/spark/job-server/spark_job_server.gradle index 7e2deaf6e395..42691461c3ef 100644 --- a/runners/spark/job-server/spark_job_server.gradle +++ b/runners/spark/job-server/spark_job_server.gradle @@ -28,7 +28,10 @@ apply plugin: 'application' // we need to set mainClassName before applying shadow plugin mainClassName = "org.apache.beam.runners.spark.SparkJobServerDriver" +def parentSparkVersion = project.parent.findProperty('spark_version') ?: '' + applyJavaNature( + requireJavaVersion: (parentSparkVersion.startsWith("4") ? org.gradle.api.JavaVersion.VERSION_17 : null), automaticModuleName: 'org.apache.beam.runners.spark.jobserver', archivesBaseName: project.hasProperty('archives_base_name') ? archives_base_name : archivesBaseName, validateShadowJar: false, diff --git a/runners/spark/spark_runner.gradle b/runners/spark/spark_runner.gradle index 091cfb053f2e..312a81cc927b 100644 --- a/runners/spark/spark_runner.gradle +++ b/runners/spark/spark_runner.gradle @@ -21,6 +21,7 @@ import groovy.json.JsonOutput apply plugin: 'org.apache.beam.module' applyJavaNature( enableStrictDependencies: true, + requireJavaVersion: (spark_version.startsWith("4") ? org.gradle.api.JavaVersion.VERSION_17 : null), automaticModuleName: 'org.apache.beam.runners.spark', archivesBaseName: (project.hasProperty('archives_base_name') ? archives_base_name : archivesBaseName), exportJavadoc: (project.hasProperty('exportJavadoc') ? exportJavadoc : true), @@ -35,6 +36,16 @@ applyJavaNature( description = "Apache Beam :: Runners :: Spark $spark_version" +// Numeric version comparison (lexicographic string compare was fragile — e.g. "3.10.0" < "3.5.0"). +def isSparkAtLeast = { String minVersion -> + def parts = spark_version.tokenize('.-').findAll { it.isInteger() }*.toInteger() + def minParts = minVersion.tokenize('.')*.toInteger() + for (int i = 0; i < Math.min(parts.size(), minParts.size()); i++) { + if (parts[i] != minParts[i]) return parts[i] > minParts[i] + } + return parts.size() >= minParts.size() +} + /* * We need to rely on manually specifying these evaluationDependsOn to ensure that * the following projects are evaluated before we evaluate this project. This is because @@ -240,7 +251,7 @@ dependencies { spark.components.each { component -> provided "$component:$spark_version" } - if ("$spark_version" >= "3.5.0") { + if (isSparkAtLeast("3.5.0")) { implementation "org.apache.spark:spark-common-utils_$spark_scala_version:$spark_version" implementation "org.apache.spark:spark-sql-api_$spark_scala_version:$spark_version" } @@ -270,7 +281,7 @@ dependencies { testImplementation library.java.mockito_core testImplementation "org.assertj:assertj-core:3.11.1" testImplementation "org.apache.zookeeper:zookeeper:3.4.11" - if ("$spark_version" >= "3.5.0") { + if (isSparkAtLeast("3.5.0")) { testImplementation "org.apache.spark:spark-common-utils_$spark_scala_version:$spark_version" testImplementation "org.apache.spark:spark-sql-api_$spark_scala_version:$spark_version" } @@ -284,7 +295,7 @@ dependencies { "hadoopVersion$kv.key" "org.apache.hadoop:hadoop-common:$kv.value" // Force paranamer 2.8 to avoid issues when using Scala 2.12 "hadoopVersion$kv.key" "com.thoughtworks.paranamer:paranamer:2.8" - if ("$spark_version" >= "3.5.0") { + if (isSparkAtLeast("3.5.0")) { // Add log4j 2.x dependencies as Spark 3.5+ uses slf4j with log4j 2.x backend "hadoopVersion$kv.key" library.java.log4j2_api "hadoopVersion$kv.key" library.java.log4j2_core @@ -310,7 +321,7 @@ configurations.validatesRunner { // Exclude to make sure log4j binding is used exclude group: "org.slf4j", module: "slf4j-simple" - if ("$spark_version" >= "3.5.0") { + if (isSparkAtLeast("3.5.0")) { // Exclude log4j 1.x dependencies to prevent conflict with log4j 2.x used by spark 3.5+ exclude group: "log4j", module: "log4j" } @@ -321,7 +332,7 @@ hadoopVersions.each { kv -> resolutionStrategy { force "org.apache.hadoop:hadoop-common:$kv.value" } - if ("$spark_version" >= "3.5.0") { + if (isSparkAtLeast("3.5.0")) { // Exclude log4j 1.x dependencies to prevent conflict with log4j 2.x used by spark 3.5+ exclude group: "log4j", module: "log4j" } From 697843323c2ce5bd431cb6b84318be6e0940ac1c Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 16:18:27 +0000 Subject: [PATCH 4/8] CHANGES.md: link Spark 4 prep entry to PR #38324 --- CHANGES.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGES.md b/CHANGES.md index d80b6affbb92..346313e5c8c1 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -81,7 +81,7 @@ Spark 4, introduced a numeric `isSparkAtLeast` Gradle helper to replace lexicographic Spark version comparison, and routed `requireJavaVersion` through `applyJavaNature` so future Spark 4 builds enforce Java 17. No - behavior change on Spark 3.5. + behavior change on Spark 3.5 ([#38324](https://github.com/apache/beam/pull/38324)). ## Breaking Changes From 41f9df0ff9d00b99d0c681f674d3f4525a44b98d Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 16:24:02 +0000 Subject: [PATCH 5/8] [Spark Runner] Address Gemini review on #38324 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - spark_runner.gradle: move isSparkAtLeast closure above applyJavaNature and use isSparkAtLeast("4.0.0") for the requireJavaVersion gate instead of spark_version.startsWith("4"), for consistency and robustness against future major version strings (e.g. "4.0.0-preview"). - spark_runner.gradle: parse minVersion the same way as spark_version inside isSparkAtLeast — tokenize('.-').findAll { it.isInteger() } — so a callsite using a suffixed version string can't trigger NumberFormatException. - spark_job_server.gradle: replace project.parent.findProperty (NPE if the project has no parent) with project.findProperty, which already walks up the project hierarchy. Renamed the local var to sparkVersion to match the inheritance semantics. --- .../spark/job-server/spark_job_server.gradle | 4 ++-- runners/spark/spark_runner.gradle | 23 ++++++++++--------- 2 files changed, 14 insertions(+), 13 deletions(-) diff --git a/runners/spark/job-server/spark_job_server.gradle b/runners/spark/job-server/spark_job_server.gradle index 42691461c3ef..a35bbea377d5 100644 --- a/runners/spark/job-server/spark_job_server.gradle +++ b/runners/spark/job-server/spark_job_server.gradle @@ -28,10 +28,10 @@ apply plugin: 'application' // we need to set mainClassName before applying shadow plugin mainClassName = "org.apache.beam.runners.spark.SparkJobServerDriver" -def parentSparkVersion = project.parent.findProperty('spark_version') ?: '' +def sparkVersion = project.findProperty('spark_version') ?: '' applyJavaNature( - requireJavaVersion: (parentSparkVersion.startsWith("4") ? org.gradle.api.JavaVersion.VERSION_17 : null), + requireJavaVersion: (sparkVersion.startsWith("4") ? org.gradle.api.JavaVersion.VERSION_17 : null), automaticModuleName: 'org.apache.beam.runners.spark.jobserver', archivesBaseName: project.hasProperty('archives_base_name') ? archives_base_name : archivesBaseName, validateShadowJar: false, diff --git a/runners/spark/spark_runner.gradle b/runners/spark/spark_runner.gradle index 312a81cc927b..0e77821e533e 100644 --- a/runners/spark/spark_runner.gradle +++ b/runners/spark/spark_runner.gradle @@ -19,9 +19,20 @@ import groovy.json.JsonOutput apply plugin: 'org.apache.beam.module' + +// Numeric version comparison (lexicographic string compare was fragile — e.g. "3.10.0" < "3.5.0"). +def isSparkAtLeast = { String minVersion -> + def parts = spark_version.tokenize('.-').findAll { it.isInteger() }*.toInteger() + def minParts = minVersion.tokenize('.-').findAll { it.isInteger() }*.toInteger() + for (int i = 0; i < Math.min(parts.size(), minParts.size()); i++) { + if (parts[i] != minParts[i]) return parts[i] > minParts[i] + } + return parts.size() >= minParts.size() +} + applyJavaNature( enableStrictDependencies: true, - requireJavaVersion: (spark_version.startsWith("4") ? org.gradle.api.JavaVersion.VERSION_17 : null), + requireJavaVersion: (isSparkAtLeast("4.0.0") ? org.gradle.api.JavaVersion.VERSION_17 : null), automaticModuleName: 'org.apache.beam.runners.spark', archivesBaseName: (project.hasProperty('archives_base_name') ? archives_base_name : archivesBaseName), exportJavadoc: (project.hasProperty('exportJavadoc') ? exportJavadoc : true), @@ -36,16 +47,6 @@ applyJavaNature( description = "Apache Beam :: Runners :: Spark $spark_version" -// Numeric version comparison (lexicographic string compare was fragile — e.g. "3.10.0" < "3.5.0"). -def isSparkAtLeast = { String minVersion -> - def parts = spark_version.tokenize('.-').findAll { it.isInteger() }*.toInteger() - def minParts = minVersion.tokenize('.')*.toInteger() - for (int i = 0; i < Math.min(parts.size(), minParts.size()); i++) { - if (parts[i] != minParts[i]) return parts[i] > minParts[i] - } - return parts.size() >= minParts.size() -} - /* * We need to rely on manually specifying these evaluationDependsOn to ensure that * the following projects are evaluated before we evaluate this project. This is because From 07be39e1ca960b4ae903311b4f7fe8886be29287 Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 17:05:42 +0000 Subject: [PATCH 6/8] CHANGES.md: use /issues/ URL form for #38324 link (validator) --- CHANGES.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGES.md b/CHANGES.md index 346313e5c8c1..6b554dec9704 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -81,7 +81,7 @@ Spark 4, introduced a numeric `isSparkAtLeast` Gradle helper to replace lexicographic Spark version comparison, and routed `requireJavaVersion` through `applyJavaNature` so future Spark 4 builds enforce Java 17. No - behavior change on Spark 3.5 ([#38324](https://github.com/apache/beam/pull/38324)). + behavior change on Spark 3.5 ([#38324](https://github.com/apache/beam/issues/38324)). ## Breaking Changes From 075b3a4e6a01fb3f3bca4b467f07b472a375fcea Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Wed, 29 Apr 2026 21:54:12 +0200 Subject: [PATCH 7/8] Clear CHANGES.MD --- CHANGES.md | 6 ------ 1 file changed, 6 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index 6b554dec9704..e8e11e830d14 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -76,12 +76,6 @@ ([#38139](https://github.com/apache/beam/issues/38139)). * (Python) Added type alias for with_exception_handling to be used for typehints. ([#38173](https://github.com/apache/beam/issues/38173)). * Added plugin mechanism to support different Lineage implementations (Java) ([#36790](https://github.com/apache/beam/issues/36790)). -* Prepared the shared Spark runner base for Spark 4 compatibility: migrated - Scala collection and shaded-Guava calls to forms valid on both Spark 3 and - Spark 4, introduced a numeric `isSparkAtLeast` Gradle helper to replace - lexicographic Spark version comparison, and routed `requireJavaVersion` - through `applyJavaNature` so future Spark 4 builds enforce Java 17. No - behavior change on Spark 3.5 ([#38324](https://github.com/apache/beam/issues/38324)). ## Breaking Changes From a076c4f61b26df0864f26b8ebeeda3a84f7fb262 Mon Sep 17 00:00:00 2001 From: Tobias Kaymak Date: Thu, 30 Apr 2026 07:11:10 +0000 Subject: [PATCH 8/8] flaky FileIOTest #19480 + KinesisIOWriteTest retry