This is an automated email from the ASF dual-hosted git repository.
Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new d32207058b0 add Spark job-server --add-opens (#39339)
d32207058b0 is described below
commit d32207058b0554411b45231280d26fcb0355d459
Author: Abdelrahman Ibrahim <[email protected]>
AuthorDate: Wed Jul 15 15:28:53 2026 +0300
add Spark job-server --add-opens (#39339)
---
.github/trigger_files/beam_PostCommit_Go_VR_Spark.json | 4 ++--
.../trigger_files/beam_PostCommit_Python_Examples_Spark.json | 4 ++++
.../beam_PostCommit_Python_ValidatesRunner_Spark.json | 3 ++-
sdks/go/test/run_validatesrunner_tests.sh | 4 ++++
sdks/python/apache_beam/runners/portability/job_server.py | 4 ++--
.../python/apache_beam/runners/portability/job_server_test.py | 2 +-
sdks/python/apache_beam/runners/portability/spark_runner.py | 11 +++++++++++
.../apache_beam/runners/portability/spark_runner_test.py | 2 ++
8 files changed, 28 insertions(+), 6 deletions(-)
diff --git a/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json
b/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json
index 72b690e649d..0dcf05b4215 100644
--- a/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json
+++ b/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json
@@ -1,5 +1,5 @@
{
"comment": "Modify this file in a trivial way to cause this test suite to
run",
- "modification": 1,
- "https://github.com/apache/beam/pull/36527": "skip a processing time timer
test in spark",
+ "modification": 2,
+ "https://github.com/apache/beam/pull/36527": "skip a processing time timer
test in spark"
}
diff --git a/.github/trigger_files/beam_PostCommit_Python_Examples_Spark.json
b/.github/trigger_files/beam_PostCommit_Python_Examples_Spark.json
new file mode 100644
index 00000000000..e3d6056a5de
--- /dev/null
+++ b/.github/trigger_files/beam_PostCommit_Python_Examples_Spark.json
@@ -0,0 +1,4 @@
+{
+ "comment": "Modify this file in a trivial way to cause this test suite to
run",
+ "modification": 1
+}
diff --git
a/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json
b/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json
index f4288649d8a..f4ec72dc416 100644
--- a/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json
+++ b/.github/trigger_files/beam_PostCommit_Python_ValidatesRunner_Spark.json
@@ -2,5 +2,6 @@
"https://github.com/apache/beam/pull/34830": "testing",
"https://github.com/apache/beam/issues/35429": "testing",
"trigger-2026-04-04": "portable_runner expand_sdf opt-in",
- "https://github.com/apache/beam/pull/38892": "UnboundedSource portable VR
test"
+ "https://github.com/apache/beam/pull/38892": "UnboundedSource portable VR
test",
+ "modification": 1
}
diff --git a/sdks/go/test/run_validatesrunner_tests.sh
b/sdks/go/test/run_validatesrunner_tests.sh
index 9ea89fd22cf..f1d68126769 100755
--- a/sdks/go/test/run_validatesrunner_tests.sh
+++ b/sdks/go/test/run_validatesrunner_tests.sh
@@ -286,6 +286,10 @@ if [[ "$RUNNER" == "flink" || "$RUNNER" == "spark" ||
"$RUNNER" == "portable" ||
--artifact-port 0 &
elif [[ "$RUNNER" == "spark" ]]; then
"$JAVA_CMD" \
+ --add-opens=java.base/sun.nio.ch=ALL-UNNAMED \
+ --add-opens=java.base/java.nio=ALL-UNNAMED \
+ --add-opens=java.base/java.util=ALL-UNNAMED \
+ --add-opens=java.base/java.lang.invoke=ALL-UNNAMED \
-jar $SPARK_JOB_SERVER_JAR \
--spark-master-url local \
--job-port $JOB_PORT \
diff --git a/sdks/python/apache_beam/runners/portability/job_server.py
b/sdks/python/apache_beam/runners/portability/job_server.py
index 8bdad6b70fe..6c55f642c61 100644
--- a/sdks/python/apache_beam/runners/portability/job_server.py
+++ b/sdks/python/apache_beam/runners/portability/job_server.py
@@ -175,8 +175,8 @@ class JavaJarJobServer(SubprocessJobServer):
self._artifacts_dir if self._artifacts_dir else self.local_temp_dir(
prefix='artifacts'))
job_port, = subprocess_server.pick_port(self._job_port)
- subprocess_cmd = [self._java_launcher, '-jar'] + self._jvm_properties + [
- jar_path
+ subprocess_cmd = [self._java_launcher] + list(self._jvm_properties) + [
+ '-jar', jar_path
] + list(
self.java_arguments(
job_port, self._artifact_port, self._expansion_port,
artifacts_dir))
diff --git a/sdks/python/apache_beam/runners/portability/job_server_test.py
b/sdks/python/apache_beam/runners/portability/job_server_test.py
index 13b3629b24b..fbf66855779 100644
--- a/sdks/python/apache_beam/runners/portability/job_server_test.py
+++ b/sdks/python/apache_beam/runners/portability/job_server_test.py
@@ -62,8 +62,8 @@ class JavaJarJobServerTest(unittest.TestCase):
subprocess_cmd,
[
'/path/to/java',
- '-jar',
'-Dsome.property=value',
+ '-jar',
'/path/to/jar',
'--artifacts-dir',
'/path/to/artifacts/',
diff --git a/sdks/python/apache_beam/runners/portability/spark_runner.py
b/sdks/python/apache_beam/runners/portability/spark_runner.py
index 480fbdecdce..f7eeb6509f6 100644
--- a/sdks/python/apache_beam/runners/portability/spark_runner.py
+++ b/sdks/python/apache_beam/runners/portability/spark_runner.py
@@ -31,6 +31,13 @@ from apache_beam.runners.portability import
spark_uber_jar_job_server
# https://spark.apache.org/docs/latest/submitting-applications.html#master-urls
LOCAL_MASTER_PATTERN = r'^local(\[.+\])?$'
+SPARK_JAR_JOB_SERVER_JVM_ARGS = [
+ '--add-opens=java.base/sun.nio.ch=ALL-UNNAMED',
+ '--add-opens=java.base/java.nio=ALL-UNNAMED',
+ '--add-opens=java.base/java.util=ALL-UNNAMED',
+ '--add-opens=java.base/java.lang.invoke=ALL-UNNAMED',
+]
+
# Since Java job servers are heavyweight external processes, cache them.
# This applies only to SparkJarJobServer, not SparkUberJarJobServer.
JOB_SERVER_CACHE = {}
@@ -84,6 +91,10 @@ class SparkJarJobServer(job_server.JavaJarJobServer):
self._jar = options.spark_job_server_jar
self._master_url = options.spark_master_url
self._spark_version = options.spark_version
+ self._jvm_properties = list(self._jvm_properties)
+ for arg in SPARK_JAR_JOB_SERVER_JVM_ARGS:
+ if arg not in self._jvm_properties:
+ self._jvm_properties.append(arg)
def path_to_jar(self):
if self._jar:
diff --git a/sdks/python/apache_beam/runners/portability/spark_runner_test.py
b/sdks/python/apache_beam/runners/portability/spark_runner_test.py
index 54238bbf192..4152b8d09f4 100644
--- a/sdks/python/apache_beam/runners/portability/spark_runner_test.py
+++ b/sdks/python/apache_beam/runners/portability/spark_runner_test.py
@@ -29,6 +29,7 @@ from apache_beam.options.pipeline_options import
PortableOptions
from apache_beam.runners.portability import job_server
from apache_beam.runners.portability import portable_runner
from apache_beam.runners.portability import portable_runner_test
+from apache_beam.runners.portability.spark_runner import
SPARK_JAR_JOB_SERVER_JVM_ARGS
from apache_beam.utils import subprocess_server
# Run as
@@ -100,6 +101,7 @@ class
SparkRunnerTest(portable_runner_test.PortableRunnerTest):
try:
return [
subprocess_server.JavaHelper.get_java(),
+ *SPARK_JAR_JOB_SERVER_JVM_ARGS,
'-Dbeam.spark.test.reuseSparkContext=true',
'-jar',
cls.spark_job_server_jar,