diff --git a/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json b/.github/trigger_files/beam_PostCommit_Go_VR_Spark.json index 72b690e649d3..0dcf05b4215b 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 000000000000..e3d6056a5de9 --- /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 f4288649d8ab..f4ec72dc416b 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 9ea89fd22cfc..f1d681267691 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 8bdad6b70fea..6c55f642c615 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 @@ def subprocess_cmd_and_endpoint(self): 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 13b3629b24bf..fbf668557794 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 @@ def test_subprocess_cmd_and_endpoint(self): 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 480fbdecdce3..f7eeb6509f6a 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 @@ # 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 @@ def __init__(self, options): 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 54238bbf1928..4152b8d09f4f 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.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 @@ def _subprocess_command(cls, job_port, expansion_port): try: return [ subprocess_server.JavaHelper.get_java(), + *SPARK_JAR_JOB_SERVER_JVM_ARGS, '-Dbeam.spark.test.reuseSparkContext=true', '-jar', cls.spark_job_server_jar,